diff --git a/README.md b/README.md index 2701740..8a98a37 100644 --- a/README.md +++ b/README.md @@ -47,7 +47,7 @@ Rust crate 位于 `crates/prompt-ir`,公开 JSON Schema 位于 `schemas/prompt cargo run -p codeischeap-desktop-api --bin export-desktop-contract ``` -捕获 sidecar 位于 `sidecars/mitmproxy`,跨进程契约位于 `crates/capture-ipc`,公开 CaptureEnvelope Schema 位于 `schemas/capture-envelope/v0.1.schema.json`。sidecar 的安装、测试与打包命令见 [`sidecars/mitmproxy/README.md`](./sidecars/mitmproxy/README.md)。 +捕获 sidecar 位于 `sidecars/mitmproxy`,跨进程契约位于 `crates/capture-ipc`,公开 CaptureEnvelope Schema 位于 `schemas/capture-envelope/v0.1.schema.json`。打包探针会验证 HTTP/1.1、HTTP/2、压缩、流式脱敏、客户端取消与捕获 IPC 背压;捕获队列满时只丢弃记录,不阻塞代理转发。sidecar 的安装、测试与打包命令见 [`sidecars/mitmproxy/README.md`](./sidecars/mitmproxy/README.md)。 捕获范围与敏感字段由 `policies/capture-policy.v0.1.json` 定义,公开 schema 位于 `schemas/capture-policy/v0.1.schema.json`。Python sidecar 在 IPC 前执行策略,`crates/core` 在进入持久化前再次拒绝越界请求并删除遗漏凭据。 diff --git a/crates/sidecar-runtime/src/lib.rs b/crates/sidecar-runtime/src/lib.rs index c422633..b40f759 100644 --- a/crates/sidecar-runtime/src/lib.rs +++ b/crates/sidecar-runtime/src/lib.rs @@ -1592,6 +1592,8 @@ struct IntegrationProbe { non_target_tunnel: bool, http2_preserved: bool, transport_context_preserved: bool, + client_cancellation_survived: bool, + capture_backpressure_nonblocking: bool, } fn validate_manifest( @@ -1650,6 +1652,8 @@ fn validate_manifest( || !probe.non_target_tunnel || !probe.http2_preserved || !probe.transport_context_preserved + || !probe.client_cancellation_survived + || !probe.capture_backpressure_nonblocking || probe.credential_canaries_in_envelope != 0 { return Err(SidecarError::InvalidManifest( @@ -2196,7 +2200,9 @@ MAoGCCqGSM49BAMCA0gAMEUCIQC1PB8+NumezrQf5unFGhVeufUcyw/sjH6p1aqs "stream_credentials_removed": true, "non_target_tunnel": true, "http2_preserved": true, - "transport_context_preserved": true + "transport_context_preserved": true, + "client_cancellation_survived": true, + "capture_backpressure_nonblocking": true }, "bundle_ready": true, "release_ready": signature == "valid" diff --git a/docs/progress.html b/docs/progress.html index 0334965..b81a49e 100644 --- a/docs/progress.html +++ b/docs/progress.html @@ -62,7 +62,7 @@

4. 工作流进度

Capture / NetworkCAP-001~00785%In progressCodex / TBD补齐 macOS 特权代理 helper;生产签名随发布凭据补齐 Prompt / AdaptersPAR-001~007100%DoneCodex / TBD保持价格目录与 provider usage 映射可追溯 Data / SecurityDAT-001~002、SEC-001~00475%In progressCodex / TBD继续 OS 级 IPC ACL、WASI 权限、供应链控制与独立安全评审 - Test / ReleaseTST-001~004、REL-001~00348%In progressCodex / TBD继续 TST-002 Proxy 取消/背压与支持处理流程 + Test / ReleaseTST-001~004、REL-001~00355%In progressCodex / TBD继续 TST-002 完整协议一致性与支持处理流程

5. 当前迭代:S5 / S6

@@ -98,7 +98,7 @@

5. 当前迭代:S5 / S6

PAR-006Ollama 适配器Done100%Codex / TBD2026-07-17本地 /api/chat 与 /api/generate、system/messages/images/tools/options、JSON/NDJSON 响应、usage 与 Raw 精确定位完成 PAR-007token、成本和语义指纹Done100%Codex / TBD2026-07-17四厂商 reported usage 归一化、显式 estimated 估算、版本化价格匹配、未知价格留空与 BLAKE3-256 语义指纹完成 TST-001协议 fixture 与 golden testsDone100%Codex / TBD2026-07-17版本化能力矩阵覆盖 OpenAI、Anthropic、Gemini 与 Ollama 的请求、响应、流式、工具、多模态、错误和 Raw fallback;声明均由 fixture 与 golden 验证 - TST-002Gateway/Proxy 集成测试In progress65%Codex / TBD2026-08-25认证 IPC、暂停丢弃、真实模式切换与进程树清理、压缩、流式脱敏、非目标 TLS、HTTP/2 基线一致性及 Windows CA 精确增删通过;真实 ROOT 往返需交互式 Windows 会话,Proxy 取消、背压与完整协议一致性待扩展 + TST-002Gateway/Proxy 集成测试In progress82%Codex / TBD2026-08-25认证 IPC、暂停丢弃、真实模式切换与进程树清理、压缩、流式脱敏、非目标 TLS、HTTP/2 基线一致性及 Windows CA 精确增删通过;真实 sidecar 在客户端取消后保持存活,IPC 停止消费且 72 个事件超过容量 64 时,36 个目标响应仍在 15 秒内完成;真实 ROOT 往返和其余协议一致性待扩展 REL-002诊断与支持包In progress75%Codex / Support Owner TBD2026-09-22版本化 JSON 支持包可预览、复制和保存,包含不带请求标识的兼容诊断树;256 KiB code-only journal 与最近 100 条事件接入,排除 Prompt、Raw 和日志详情;支持处理流程待完成 SPIKE-001Gateway 流式透明转发验证Done100%Codex / TBD2026-07-14双向流式、取消传递、头清理与稳定 502 集成测试通过 SPIKE-002mitmproxy sidecar IPC/打包验证Done100%Codex / TBD2026-07-14凭据清理、IPC、打包与真实转发通过 @@ -121,7 +121,7 @@

6. 后续迭代承诺

S1APP-001/002、DAT-001、SEC-001、Core 事件骨架In progress加密桌面链路与 SEC-001 已通过;Core event 与独立安全评审待完成 S2CAP-001/002、PAR-002/003、APP-003DoneGateway、OpenAI 解析与千条实时工作台全部通过验收 S3PAR-004、APP-004、TST-001、DAT-002Done双厂商 Inspector、能力矩阵与数据生命周期全部通过验收 - S4CAP-003~005、TST-002In progresssidecar bundle、桌面运行时、协议矩阵、跨平台 CA 状态及两平台用户级信任生命周期已实现;签名、Proxy 取消/背压与交互式验收待推进 + S4CAP-003~005、TST-002In progresssidecar bundle、桌面运行时、协议矩阵、Proxy 取消/背压、跨平台 CA 状态及两平台用户级信任生命周期已实现;签名、其余协议一致性与交互式验收待推进 S5CAP-006/007、SEC-002/003、APP-006、TST-003In progressSEC-002、CAP-007 完成,IPC 抗阻塞、双模式 OS PID 归因和兼容诊断树已接入;继续 macOS helper、OS socket ACL 与其余故障注入 S6PAR-005~007、APP-005Done四厂商适配器、token/成本/指纹、全文搜索和结构/文本 Compare 全部完成 S7TST-004、性能、可访问性、诊断与保留Not started功能冻结 @@ -218,6 +218,7 @@

11. 决策与变更记录

2026-07-18Capture / SecurityCAP-007 完成、SEC-003 推进IPC 0.3 在认证帧中传递严格 loopback 临时端点,桌面端忽略 sidecar PID 并用 OS socket 表精确归因 Proxy 进程;端点不进入 Envelope、数据库或导出,查询失败保持未知Codex 2026-07-18Security / IPCSEC-003 sidecar peer 绑定完成首段IPC 0.4 让 sidecar 等待服务端 ACK,从协议层保持首连接直到 Windows/macOS OS socket 表完成 PID 核对;不匹配、查询失败或 2 秒超时只拒绝捕获,代理转发继续,后续连接仍受会话 token 保护Codex 2026-07-18Security / IPCSEC-003 Linux peer 绑定完成Linux 读取当前 network namespace 的 /proc/net/tcp{,6},以客户端到 IPC 服务端的精确四元组定位 ESTABLISHED socket inode,再扫描 /proc/<pid>/fd;仅唯一 owner 且 PID 等于已启动 sidecar 时接受首帧,权限受限、连接消失、歧义和解析失败均保持拒绝Codex + 2026-07-18Test / ProxyTST-002 取消与背压矩阵完成单元测试以阻塞 IPC worker 和容量 1 队列证明满队列 submit 立即失败;首次 IPC 失败会清空 backlog、熔断 10 秒并自动探测恢复;真实打包 sidecar 在 IPC listener 停止消费后承受 36 个目标请求、72 个捕获事件和一次客户端取消,转发响应保持一致且进程继续存活;两个结论成为 Python 与 Rust bundle 强制合同Codex 后续范围、架构、日期或资源变化均在此追加,并链接对应 ADR/会议结论。 diff --git a/sidecars/mitmproxy/README.md b/sidecars/mitmproxy/README.md index 41a9b4b..818b3ed 100644 --- a/sidecars/mitmproxy/README.md +++ b/sidecars/mitmproxy/README.md @@ -12,6 +12,8 @@ CIC_CAPTURE_HOSTS=api.openai.com,api.anthropic.com The addon loads `policies/capture-policy.v0.1.json` and only records exact target hosts, approved POST paths, and methods. `CIC_CAPTURE_HOSTS` can opt an OpenAI-compatible host into the same approved path set; it does not disable path checks. The addon removes credential headers, sensitive query fields, and recursively named JSON secret fields before sending an authenticated NDJSON envelope. Unsupported request body formats are omitted. The original network request is not modified. +Capture delivery uses a bounded non-blocking queue. A full queue drops only the capture event. The first IPC delivery failure discards the pending backlog and opens a 10-second retry circuit; proxy forwarding continues, and the next event after the cooldown probes IPC recovery. + ```powershell python -m pip install -r sidecars/mitmproxy/requirements-build.txt python -m unittest discover -s sidecars/mitmproxy/tests -v diff --git a/sidecars/mitmproxy/codeischeap_addon.py b/sidecars/mitmproxy/codeischeap_addon.py index 53b2f83..04a0475 100644 --- a/sidecars/mitmproxy/codeischeap_addon.py +++ b/sidecars/mitmproxy/codeischeap_addon.py @@ -550,23 +550,49 @@ def send_envelope( ).encode("utf-8") with socket.create_connection((config.host, config.port), timeout=config.timeout_seconds) as connection: connection.sendall(frames) - acknowledgement = connection.makefile("rb").readline(256) - if json.loads(acknowledgement) != {"status": "accepted"}: + acknowledgement = bytearray() + while len(acknowledgement) < 256 and b"\n" not in acknowledgement: + chunk = connection.recv(256 - len(acknowledgement)) + if not chunk: + break + acknowledgement.extend(chunk) + frame, separator, _ = bytes(acknowledgement).partition(b"\n") + if not separator or json.loads(frame) != {"status": "accepted"}: raise ValueError("capture IPC acknowledgement is invalid") class IpcEmitter: - def __init__(self, config: IpcConfig, capacity: int = 64) -> None: + def __init__( + self, + config: IpcConfig, + capacity: int = 64, + retry_delay_seconds: float = 10.0, + ) -> None: self._config = config + self._retry_delay_seconds = retry_delay_seconds self._queue: queue.Queue[ tuple[dict[str, Any], dict[str, str] | None] | None ] = queue.Queue(maxsize=capacity) + self._state_lock = threading.Lock() + self._retry_after = 0.0 + self._delivery_unavailable = threading.Event() self._worker = threading.Thread(target=self._run, name="codeischeap-ipc", daemon=True) self._worker.start() def submit( self, envelope: dict[str, Any], transport: dict[str, str] | None ) -> bool: + now = time.monotonic() + with self._state_lock: + if now < self._retry_after: + unavailable = True + else: + if self._retry_after: + self._retry_after = 0.0 + self._delivery_unavailable.clear() + unavailable = False + if unavailable: + return False try: self._queue.put_nowait((envelope, transport)) return True @@ -589,7 +615,20 @@ def _run(self) -> None: try: send_envelope(self._config, envelope, transport) except (OSError, ValueError): - _warn("CodeIsCheap capture IPC delivery failed; the request was not recorded") + with self._state_lock: + self._retry_after = time.monotonic() + self._retry_delay_seconds + self._delivery_unavailable.set() + if self._discard_pending(): + return + + def _discard_pending(self) -> bool: + while True: + try: + submission = self._queue.get_nowait() + except queue.Empty: + return False + if submission is None: + return True def _warn(message: str) -> None: @@ -597,7 +636,7 @@ def _warn(message: str) -> None: from mitmproxy import ctx ctx.log.warn(message) - except (ImportError, RuntimeError): + except (AttributeError, ImportError, RuntimeError): pass @@ -638,8 +677,7 @@ def _submit( except (AttributeError, TypeError, ValueError): _warn(failure_message) return - if not self._emitter.submit(envelope, build_transport_context(flow)): - _warn("CodeIsCheap capture queue is full; the request was not recorded") + self._emitter.submit(envelope, build_transport_context(flow)) def request(self, flow: Any) -> None: self._submit( diff --git a/sidecars/mitmproxy/package_sidecar.py b/sidecars/mitmproxy/package_sidecar.py index aa27d28..e47fa13 100644 --- a/sidecars/mitmproxy/package_sidecar.py +++ b/sidecars/mitmproxy/package_sidecar.py @@ -299,6 +299,8 @@ def main() -> None: "non_target_tunnel", "http2_preserved", "transport_context_preserved", + "client_cancellation_survived", + "capture_backpressure_nonblocking", ) ) and probe_result.get("credential_canaries_in_envelope") == 0 manifest = { diff --git a/sidecars/mitmproxy/tests/test_addon.py b/sidecars/mitmproxy/tests/test_addon.py index eb912bf..faf8a3d 100644 --- a/sidecars/mitmproxy/tests/test_addon.py +++ b/sidecars/mitmproxy/tests/test_addon.py @@ -7,6 +7,7 @@ import sys import tempfile import threading +import time import unittest from unittest.mock import patch @@ -17,6 +18,7 @@ sys.path.insert(0, str(SIDECAR_DIR)) from codeischeap_addon import ( + IpcEmitter, IpcConfig, build_envelope, build_failure_envelope, @@ -398,7 +400,9 @@ def accept() -> None: with connection: stream = connection.makefile("rb") received.extend([stream.readline(), stream.readline()]) - connection.sendall(b'{"status":"accepted"}\n') + connection.sendall(b'{"status":') + time.sleep(0.01) + connection.sendall(b'"accepted"}\n') worker = threading.Thread(target=accept) worker.start() @@ -420,6 +424,90 @@ def accept() -> None: self.assertNotIn("transport", captured) self.assertEqual(captured, envelope) + def test_ipc_queue_backpressure_never_blocks_flow_submission(self) -> None: + first_delivery_started = threading.Event() + release_first_delivery = threading.Event() + second_delivery_finished = threading.Event() + deliveries: list[dict[str, object]] = [] + + def blocked_send( + config: IpcConfig, + envelope: dict[str, object], + transport: dict[str, str] | None = None, + ) -> None: + del config, transport + deliveries.append(envelope) + if len(deliveries) == 1: + first_delivery_started.set() + release_first_delivery.wait(2) + elif len(deliveries) == 2: + second_delivery_finished.set() + + config = IpcConfig("127.0.0.1", 1, "synthetic-token") + with patch("codeischeap_addon.send_envelope", side_effect=blocked_send): + emitter = IpcEmitter(config, capacity=1) + try: + self.assertTrue(emitter.submit({"capture_id": "first"}, None)) + self.assertTrue(first_delivery_started.wait(1)) + self.assertTrue(emitter.submit({"capture_id": "queued"}, None)) + + started = time.monotonic() + self.assertFalse(emitter.submit({"capture_id": "dropped"}, None)) + self.assertLess(time.monotonic() - started, 0.25) + finally: + release_first_delivery.set() + + self.assertTrue(second_delivery_finished.wait(1)) + emitter.close() + self.assertEqual( + [delivery["capture_id"] for delivery in deliveries], + ["first", "queued"], + ) + + def test_ipc_failure_discards_backlog_and_opens_a_retry_circuit(self) -> None: + delivery_attempted = threading.Event() + release_failure = threading.Event() + recovery_delivered = threading.Event() + attempts: list[str] = [] + + def failed_send( + config: IpcConfig, + envelope: dict[str, object], + transport: dict[str, str] | None = None, + ) -> None: + del config, transport + attempts.append(str(envelope["capture_id"])) + delivery_attempted.set() + if envelope["capture_id"] == "failed": + release_failure.wait(1) + raise OSError("synthetic IPC failure") + recovery_delivered.set() + + config = IpcConfig("127.0.0.1", 1, "synthetic-token") + with patch("codeischeap_addon.send_envelope", side_effect=failed_send): + emitter = IpcEmitter(config, capacity=4, retry_delay_seconds=0.05) + try: + self.assertTrue(emitter.submit({"capture_id": "failed"}, None)) + self.assertTrue(delivery_attempted.wait(1)) + self.assertTrue(emitter.submit({"capture_id": "queued"}, None)) + release_failure.set() + self.assertTrue(emitter._delivery_unavailable.wait(1)) + + started = time.monotonic() + self.assertFalse(emitter.submit({"capture_id": "circuit-open"}, None)) + self.assertLess(time.monotonic() - started, 0.25) + + retry_deadline = time.monotonic() + 1 + while not emitter.submit({"capture_id": "recovered"}, None): + self.assertLess(time.monotonic(), retry_deadline) + time.sleep(0.01) + self.assertTrue(recovery_delivered.wait(1)) + finally: + release_failure.set() + emitter.close() + + self.assertEqual(attempts, ["failed", "recovered"]) + if __name__ == "__main__": unittest.main() diff --git a/sidecars/mitmproxy/tests/test_packaging.py b/sidecars/mitmproxy/tests/test_packaging.py index 407ab85..f25e8c4 100644 --- a/sidecars/mitmproxy/tests/test_packaging.py +++ b/sidecars/mitmproxy/tests/test_packaging.py @@ -69,6 +69,8 @@ def write_bundle(bundle: Path) -> dict: "non_target_tunnel": True, "http2_preserved": True, "transport_context_preserved": True, + "client_cancellation_survived": True, + "capture_backpressure_nonblocking": True, }, } (bundle / "sidecar-manifest.json").write_text(json.dumps(manifest)) @@ -103,12 +105,12 @@ def test_bundle_validation_checks_hash_contract_and_signature_gate(self) -> None validate_bundle(bundle) manifest["capture_contract"]["ipc_protocol"] = "0.4" - del manifest["integration_probe"]["http2_preserved"] + del manifest["integration_probe"]["capture_backpressure_nonblocking"] (bundle / "sidecar-manifest.json").write_text(json.dumps(manifest)) with self.assertRaisesRegex(ValueError, "integration probe did not pass"): validate_bundle(bundle) - manifest["integration_probe"]["http2_preserved"] = True + manifest["integration_probe"]["capture_backpressure_nonblocking"] = True (bundle / "sidecar-manifest.json").write_text(json.dumps(manifest)) with self.assertRaisesRegex(ValueError, "valid platform signature"): validate_bundle(bundle, require_signature=True) diff --git a/sidecars/mitmproxy/verify_packaged_sidecar.py b/sidecars/mitmproxy/verify_packaged_sidecar.py index c918153..f26476d 100644 --- a/sidecars/mitmproxy/verify_packaged_sidecar.py +++ b/sidecars/mitmproxy/verify_packaged_sidecar.py @@ -3,6 +3,7 @@ from __future__ import annotations import argparse +from concurrent.futures import ThreadPoolExecutor from datetime import datetime, timedelta, timezone import gzip from http.client import HTTPConnection @@ -70,17 +71,22 @@ def valid_loopback_transport(auth: dict[str, Any]) -> bool: HTTP1_CASES = ("gzip", "brotli", "sse") TARGET_CASES = (*HTTP1_CASES, "http2") +CANCEL_CASE = "cancel" +PRESSURE_CASE = "pressure" +BACKPRESSURE_REQUEST_COUNT = 36 +BACKPRESSURE_MAX_SECONDS = 15 class UpstreamHandler(BaseHTTPRequestHandler): received: dict[str, dict[str, Any]] = {} received_event = threading.Event() received_lock = threading.Lock() + cancellation_received = threading.Event() def do_POST(self) -> None: length = int(self.headers.get("content-length", "0")) case = parse_qs(urlsplit(self.path).query).get("case", [""])[0] - if case not in HTTP1_CASES: + if case not in (*HTTP1_CASES, CANCEL_CASE, PRESSURE_CASE): self.send_error(400) return received = { @@ -89,10 +95,24 @@ def do_POST(self) -> None: "x_api_key": self.headers.get("x-api-key"), "body": self.rfile.read(length).decode("utf-8"), } - with type(self).received_lock: - type(self).received[case] = received - if len(type(self).received) == len(HTTP1_CASES): - type(self).received_event.set() + if case in HTTP1_CASES: + with type(self).received_lock: + type(self).received[case] = received + if len(type(self).received) == len(HTTP1_CASES): + type(self).received_event.set() + elif case == CANCEL_CASE: + type(self).cancellation_received.set() + time.sleep(0.5) + try: + payload = b'{"ok":true,"case":"cancel"}' + self.send_response(200) + self.send_header("content-type", "application/json") + self.send_header("content-length", str(len(payload))) + self.end_headers() + self.wfile.write(payload) + except (BrokenPipeError, ConnectionAbortedError, ConnectionResetError): + pass + return self.send_response(200) if case == "sse": payload = ( @@ -114,9 +134,11 @@ def do_POST(self) -> None: if case == "gzip": payload = gzip.compress(payload) content_encoding = "gzip" - else: + elif case == "brotli": payload = brotli.compress(payload) content_encoding = "br" + else: + content_encoding = None self.send_header("content-type", content_type) if content_encoding: self.send_header("content-encoding", content_encoding) @@ -416,6 +438,32 @@ def request_target_case(proxy_port: int, upstream_port: int, case: str) -> dict[ return result +def cancel_target_request(proxy_port: int, upstream_port: int) -> None: + body = json.dumps( + {"messages": [{"role": "user", "content": "cancelled prompt"}]} + ).encode() + target = ( + f"http://localhost:{upstream_port}/v1/chat/completions?case={CANCEL_CASE}" + ) + connection = socket.create_connection(("127.0.0.1", proxy_port), timeout=10) + try: + connection.sendall( + ( + f"POST {target} HTTP/1.1\r\n" + f"Host: localhost:{upstream_port}\r\n" + "Content-Type: application/json\r\n" + f"Content-Length: {len(body)}\r\n" + "Connection: close\r\n\r\n" + ).encode() + + body + ) + if not UpstreamHandler.cancellation_received.wait(5): + raise TimeoutError("cancelled client request did not reach the fake provider") + connection.shutdown(socket.SHUT_RDWR) + finally: + connection.close() + + def open_proxy_tunnel(proxy_port: int, authority: str) -> socket.socket: connection = socket.create_connection(("127.0.0.1", proxy_port), timeout=10) connection.sendall( @@ -558,6 +606,9 @@ def main() -> None: ipc_listener.bind(("127.0.0.1", 0)) ipc_listener.listen(len(TARGET_CASES) * 2) ipc_frames: list[tuple[bytes, bytes]] = [] + stalled_ipc_connections: list[socket.socket] = [] + stop_stalled_ipc = threading.Event() + stalled_ipc_thread: threading.Thread | None = None def receive_ipc() -> None: for _ in range(len(TARGET_CASES) * 2): @@ -624,6 +675,21 @@ def receive_ipc() -> None: if ipc_thread.is_alive() or len(ipc_frames) != len(TARGET_CASES) * 2: raise TimeoutError("addon did not deliver every request and response envelope") + ipc_listener.settimeout(0.1) + + def stall_ipc() -> None: + while not stop_stalled_ipc.is_set(): + try: + connection, _ = ipc_listener.accept() + except TimeoutError: + continue + except OSError: + return + stalled_ipc_connections.append(connection) + + stalled_ipc_thread = threading.Thread(target=stall_ipc, daemon=True) + stalled_ipc_thread.start() + frames = [(json.loads(auth), json.loads(envelope)) for auth, envelope in ipc_frames] if any(auth.get("token") != token for auth, _ in frames): raise RuntimeError("IPC auth token was not preserved in every auth frame") @@ -704,6 +770,43 @@ def receive_ipc() -> None: if CANARIES["response_body"] not in forwarded_body: raise RuntimeError(f"{case} response body changed during forwarding") + pressure_started = time.monotonic() + with ThreadPoolExecutor(max_workers=12) as executor: + pressure_responses = list( + executor.map( + lambda _: request_target_case( + proxy_port, upstream.server_port, PRESSURE_CASE + ), + range(BACKPRESSURE_REQUEST_COUNT), + ) + ) + pressure_elapsed = time.monotonic() - pressure_started + if pressure_elapsed > BACKPRESSURE_MAX_SECONDS: + raise TimeoutError( + "proxy forwarding blocked behind capture IPC backpressure: " + f"{pressure_elapsed:.2f}s across " + f"{len(stalled_ipc_connections)} IPC connections" + ) + if any(response["status"] != 200 for response in pressure_responses): + raise RuntimeError("proxy failed a request while capture IPC was backpressured") + if any( + json.loads(response["body"]).get("case") != PRESSURE_CASE + for response in pressure_responses + ): + raise RuntimeError("proxy changed a response while capture IPC was backpressured") + if process.poll() is not None: + raise RuntimeError("sidecar exited after capture backpressure") + if not stalled_ipc_connections: + raise RuntimeError("capture IPC backpressure connection was not established") + stop_stalled_ipc.set() + stalled_ipc_thread.join(timeout=1) + if stalled_ipc_thread.is_alive(): + raise TimeoutError("capture IPC backpressure listener did not stop") + for connection in stalled_ipc_connections: + connection.close() + + cancel_target_request(proxy_port, upstream.server_port) + peer_certificate, tunnel_response = request_non_target_tunnel( proxy_port, tunnel.server_port ) @@ -713,6 +816,8 @@ def receive_ipc() -> None: raise RuntimeError("non-target TLS response changed during tunneling") if not TunnelHandler.received_event.wait(5): raise TimeoutError("non-target TLS server did not receive the tunneled request") + if process.poll() is not None: + raise RuntimeError("sidecar exited after client cancellation") ipc_listener.settimeout(0.5) try: unexpected, _ = ipc_listener.accept() @@ -736,10 +841,17 @@ def receive_ipc() -> None: "stream_credentials_removed": True, "non_target_tunnel": True, "http2_preserved": True, + "client_cancellation_survived": True, + "capture_backpressure_nonblocking": True, } ) ) finally: + stop_stalled_ipc.set() + if stalled_ipc_thread is not None: + stalled_ipc_thread.join(timeout=1) + for connection in stalled_ipc_connections: + connection.close() stop_process_tree(process) upstream.shutdown() upstream.server_close() diff --git a/sidecars/mitmproxy/verify_sidecar_bundle.py b/sidecars/mitmproxy/verify_sidecar_bundle.py index 70041c1..2da2d93 100644 --- a/sidecars/mitmproxy/verify_sidecar_bundle.py +++ b/sidecars/mitmproxy/verify_sidecar_bundle.py @@ -94,6 +94,8 @@ def validate_bundle(bundle: Path, require_signature: bool = False) -> dict[str, "non_target_tunnel", "http2_preserved", "transport_context_preserved", + "client_cancellation_survived", + "capture_backpressure_nonblocking", ) ): raise ValueError("sidecar integration probe did not pass")