核心是字节流事件驱动透传与按需解码:先读初始数据识别连接类型(connect/tls/sse/http),再分发至对应处理路径;pipe中嵌入异步解码钩子,动态适配编码策略,配合流控与异常隔离保障稳定性。

核心在于不解析 HTTP 语义,而是以字节流视角做事件驱动的透传与条件解码,关键不是“代理转发”,而是“在 TCP 连接生命周期内按需触发解码逻辑”。
识别连接类型并分发处理路径
HTTP 代理必须在连接建立初期就判断流量性质,不能等完整请求到达——因为 TLS 握手、SSE 流、WebSocket Upgrade 都发生在首几个字节。需用 asyncio.StreamReader.read(1024) 读取初始数据块,再按规则分流:
- 若前 8 字节匹配
b"CONNECT "→ 启动 CONNECT 处理:解析 host:port,建立上游 TCP 连接,返回HTTP/1.1 200 Connection established,然后启动双向 pipe - 若首行含
GET|POST|HEAD且末尾为b"HTTP/1."→ 当作普通 HTTP 请求,保留原始 headers(包括 Host、Connection、Upgrade),仅对 body 做条件解码 - 若前 4 字节是
b"\x16\x03\x01\x02"(TLS ClientHello)或前 2 字节是b"\x80\x8a"(WebSocket frame)→ 直接 raw 转发,跳过所有 HTTP 解析和 header 重写 - 若首行为
b"GET /events HTTP/1.1"且含b"text/event-stream"→ 启用 SSE-aware pipe:检测data:行边界,允许按 event 分块注入清洗逻辑(如脱敏 user_id 字段)
在 pipe 中嵌入异步解码钩子
标准 pipe(reader, writer) 是纯字节搬运,要实现“流式清洗”,需把解码逻辑注册为可插拔的 handler,在数据到达时异步调用:
- 定义
async def decode_chunk(data: bytes) -> bytes:,支持 JSON body 提取、base64 解码、字段正则替换等操作 - 在 pipe 循环中,对非 TLS/raw 流量,先 await decode_chunk(data),再 write;对已知加密或二进制流(如 Protobuf over HTTP),跳过 decode 直接透传
- 解码失败时不中断连接,记录 warning 并原样转发,避免因单条脏数据导致整个流断开
- 利用
asyncio.create_task()并发执行 decode_chunk,防止 CPU 密集型解码阻塞事件循环
动态切换编码策略与上下文感知
同一个连接可能混合多种编码(如 HTTP/1.1 + chunked + gzip + JSON),需根据响应头实时调整:
- 收到 upstream 响应头后,检查
Content-Encoding: gzip或Transfer-Encoding: chunked,动态启用对应的 async decompressor(如aiohttp.ClientResponse.content的解压逻辑) - 对
Content-Type: application/json的响应,启动 JSON 流式 parser(如ijson.parse_coro),边收边提取 key 路径(如$.user.id),触发脱敏回调 - 维护 per-connection context(如 client IP、User-Agent、请求 path),用于路由不同清洗规则:移动端请求走轻量过滤,管理后台请求启用全字段审计
- 通过
asyncio.Task.current_task().get_name()标记任务来源,便于日志追踪和限流统计
保障流控与异常隔离
异步解码可能引入延迟或内存压力,必须设置硬性边界:
- 每个 pipe task 设置
asyncio.wait_for(..., timeout=30),超时强制关闭该方向连接,防止 hang 住整个 socket - 限制单次 read size ≤ 65536,避免大 payload 占满内存;对 gzip 流额外限制解压后大小(如
decompressor.decompress(data, max_length=10*1024*1024)) - 使用
asyncio.Queue(maxsize=100)缓冲 decode 后的数据,当队列满时暂停上游 read,实现反压(backpressure) - 对 decode 异常(如 JSONDecodeError)、IO 错误、SSL handshake failure,统一转为
ConnectionResetError并静默关闭对应 pipe,不影响另一方向传输











