diff --git a/Cargo.lock b/Cargo.lock index 66b748ea..88414308 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4720,7 +4720,7 @@ dependencies = [ [[package]] name = "kcp-sys" version = "0.1.0" -source = "git+https://github.com/EasyTier/kcp-sys?rev=d7427c22d764deb1860a7d37acc446ed5033464c#d7427c22d764deb1860a7d37acc446ed5033464c" +source = "git+https://github.com/EasyTier/kcp-sys?rev=268533568d734ae89dc89603078da3ca522effe1#268533568d734ae89dc89603078da3ca522effe1" dependencies = [ "anyhow", "auto_impl", diff --git a/docs/kcp-control-reliability-design-2026-09-14.md b/docs/kcp-control-reliability-design-2026-09-14.md new file mode 100644 index 00000000..1f477ff0 --- /dev/null +++ b/docs/kcp-control-reliability-design-2026-09-14.md @@ -0,0 +1,470 @@ +# KCP 控制报文可靠性与旧版兼容设计 + +日期:2026-09-14。 + +状态:已按本文实现并提交;当前实现、验收结果及仍未定位的异常见 +[实现与验证记录](kcp-control-reliability-validation-2026-09-14.md)。 +本文保留设计要求,具体通过范围以验证记录为准。 + +## 1. 背景与问题边界 + +此前已修复 TCP 代理任一方向 EOF 导致双向转发提前退出、KCP accept +通知队列满后连接交接丢失、接收缓冲排空与提前 FIN 的处理,以及握手前 +心跳触发 RST 的竞态。具体版本及结果见 +[半关闭验证记录](tcp-proxy-half-close-validation-2026-09-14.md)。 + +仍有两个问题: + +1. 已证实的握手最终 ACK 丢失:源端收到 SYNACK 后认为连接建立,目标端 + 仍停留在 SynReceived;服务端先发 greeting 的业务因此超时。 +2. 一次未主动注入丢包的反向半关闭超时:发送端出现输出队列 Full,FIN + 已进入端点输出队列,但没有证据证明它何时到达对端。失败日志未记录 + 客户端已收字节数,尚不能确定缺少的是数据、EOF,还是两者。后续 + TRACE 重跑通过,不足以消除或解释原失败。 + +KCP 数据传输自身的确认重传没有覆盖外层 SYN、SYNACK、握手确认和 FIN。 +这是连接状态机的局部协议设计问题。修复需要明确控制报文的重复处理、 +完成条件和状态寿命,无需重做 KCP 数据传输或 EasyTier 路由架构。 + +本设计解决已证实的握手恢复缺口,并为 FIN 增加可靠确认。同时继续调查 +第二项异常;不能将新增 FIN 重传直接当作该异常已经修复的证明。 + +## 2. 目标、范围与不变量 + +目标: + +- 新旧节点双向正常互通,滚动升级不得要求所有节点同时升级。 +- 新节点之间,控制报文在重试预算内丢失、重复或乱序时能够恢复。 +- 连接只交接一次,应用数据不重复交付,EOF 不越过已接收数据。 +- 半关闭只关闭一个发送方向,另一方向仍可长期传输。 +- 无法恢复时有界结束并报告错误,不能用正常 EOF 掩盖失败。 +- 关闭记录、待发送控制报文和重试任务均有明确生命周期。 + +范围:kcp-sys 的报文定义、连接状态机、端点调度与清理、必要的流错误 +传播,以及 EasyTier 的依赖 pin、集成测试和验证记录。 + +不新增路由能力公告、PeerFeatureFlag 字段、配置开关或应用层请求重试; +不修改 TCP flow-key 修复、QUIC 协议、KCP 数据分段和拥塞算法。保留现有 +connect API、调用方超时和业务测试的 5 秒 socket 超时。 + +协议不变量: + +1. 模式由本连接握手确定,交付应用后不能切换。 +2. 握手状态变化和是否产生应答是两件事;可靠模式允许状态不变但应答。 +3. 已经进入应用交接路径的连接,重复控制包不能再次 accept。 +4. FIN_ACK 只确认对方 FIN,不表示自己也关闭发送方向。 +5. 只有正常收到对端 FIN 且排空接收数据,才能向应用报告正常 EOF。 +6. FIN 确认成功后,FIN 重试期限不能限制正常半关闭连接的寿命。 +7. 重复包不能延长握手、FIN 重试和最终关闭记录的固定截止时间。 + +## 3. 版本基线与兼容证据 + +| 用途 | kcp-sys revision | +| --- | --- | +| 旧协议基线 | `d7427c22d764deb1860a7d37acc446ed5033464c` | +| 本设计实施基线,已含前述局部修复 | `b37ee660fb70bb6d816fb8bbc08b140e55e7218b` | + +源码检查确认: + +- 旧 header 含三个 u32、一个 flags 字节及一个 `rsv` 保留字节,总长 + 14 字节;旧实现不读取 `rsv`。 +- 旧 SYNACK 使用新建且清零的 header,不回显 SYN 的 `rsv`。 +- 旧 FSM 收到重复 SYN 时可能构造 RST,但状态不变;外层 + `KcpConnectionState::handle_packet` 丢弃状态未变化时的输出, + 因此已检查的旧端点不会实际发送该 RST。 +- `rsv` 未被读取、状态不变时抑制输出的行为,在已检查历史 + `a5f4b4e`、`37f653c`、`0f2fea1`、`9ce5c08`、`0f0a055`、 + `71eff18`、`d7427c2` 中一致。 + +这些是源码与历史检查,不是上述每个 revision 的运行测试。正式兼容 +声明必须对应实际测试的 EasyTier 版本及其锁定依赖,不能扩展为所有 +未检查的历史版本或第三方修改版本都已验证。 + +实现基线源码: +[报文定义](https://github.com/EasyTier/kcp-sys/blob/b37ee660fb70bb6d816fb8bbc08b140e55e7218b/src/packet_def.rs)、 +[状态机](https://github.com/EasyTier/kcp-sys/blob/b37ee660fb70bb6d816fb8bbc08b140e55e7218b/src/state.rs)、 +[端点](https://github.com/EasyTier/kcp-sys/blob/b37ee660fb70bb6d816fb8bbc08b140e55e7218b/src/endpoint.rs)。 + +## 4. 报文格式与连接内协商 + +### 4.1 字段定义 + +保持 14 字节 header、连接标识字段、已有 flags 和 KCP DATA payload +布局不变。使用已有 `rsv` 字节: + +| 值 | 定义 | +| --- | --- | +| `0` | 旧协议,后文称 legacy | +| `1` | 本文可靠控制协议,后文称可靠模式 | + +新增 FIN_ACK 使用 flags 当前未分配的最高位 `0x80`。它是独立报文, +flags 必须恰好为 FIN_ACK,`rsv=1`,payload 为空。不能用 FIN|ACK +表示关闭确认,因为旧端会把 FIN 位解释为对端关闭发送方向。 + +可靠模式报文约束: + +| 报文 | flags | payload | +| --- | --- | --- | +| SYN | SYN | 原有连接元数据,重传保持完全相同 | +| SYNACK | SYN\|ACK | 空 | +| 最终握手确认 | ACK\|DATA | 空 | +| 数据及 KCP 内部确认 | ACK\|DATA | 原有 KCP payload | +| 关闭发送方向 | FIN | 空 | +| 关闭确认 | FIN_ACK (`0x80`) | 空 | +| 重置、心跳 | 原有 RST、PING、PING\|PONG | 沿用现有定义 | + +可靠模式报文均带 `rsv=1`。legacy 的后续报文仍使用现有格式和 `rsv=0`; +新源端最初提出协商的 SYN 是可能发给旧目标端的唯一保留字节扩展。 + +只有匹配连接标识、当前握手阶段和合法 flags/payload 的 SYNACK, +才可以确认模式。不能通过 PING/PONG 确认:旧端的心跳应答会复制输入 +报文,可能保留它并不理解的 `rsv`。 + +### 4.2 协商步骤 + +1. 新源端发送 SYN,`rsv=1`,进入尚未确认模式的 SynSent。 +2. 新目标端收到 SYN:`rsv=1` 时回复 SYNACK `rsv=1`,记录可靠模式 + 候选状态;`rsv=0` 时从头使用 legacy,回复 SYNACK `rsv=0`。 +3. 旧目标端忽略 `rsv`,自然回复 SYNACK `rsv=0`。新源端收到它后, + 在交付应用前确定 legacy,最终 ACK 及后续报文都使用旧模式。 +4. 新源端收到合法 SYNACK `rsv=1` 后确认可靠模式,发送最终 ACK + `rsv=1`。目标端收到匹配模式的确认后完成握手并登记待 accept。 +5. 源端继续在收到合法 SYNACK 后完成 connect,不增加第四次握手或 + 额外 RTT。最终 ACK 丢失后的恢复规则见第 5 节。 + +未识别的 SYN `rsv` 值不得被当作可靠能力,可按不支持扩展响应 legacy +SYNACK。新源端只接受它实际支持的 SYNACK 模式 `0` 或 `1`;未知值 +报告协议错误,不启用增强行为。重复 SYN 不允许更改已有连接的模式或 +连接元数据。已协商可靠模式后,模式不匹配的控制包不得触发降级。 + +节点重启或降级后的新连接重新协商;不缓存按节点永久生效的协议模式。 +不承诺跨进程重启维持既有连接。 + +### 4.3 初始 SYN 与旧版兼容 + +首个 SYN 或 SYNACK 丢失时,源端尚不可能拿到协商确认。因此 SynSent +期间允许在原截止时间内重传完全相同的 SYN,不以前置能力公告或已收到 +SYNACK 为条件。 + +在已检查的旧端点中,重复 SYN 不引出 RST,但也不会让旧目标端重发 +SYNACK。所以新源端→旧目标端仍保留旧目标端 SYNACK 丢失时的限制。 +这一行为必须由真实旧依赖运行测试确认,不能只测试重写的 legacy 模型。 + +新目标端对旧源端保持 legacy 行为,不主动开启 SYNACK 重传。 +legacy 分支继续保留“状态不变则抑制 FSM 输出”的现有行为;不得将 +可靠模式的应答修复无条件套用,否则会改变旧连接的 RST 行为。 + +## 5. 可靠握手与重复报文 + +源端在 SynSent 重传 SYN,目标端在可靠 SynReceived 重传 SYNACK。 +所有重传使用原连接 ID,不创建额外连接或重试业务请求。 + +| 当前状态 | 输入 | 动作 | +| --- | --- | --- | +| 无状态 | 合法 SYN | 保存模式与连接元数据,建立 SynReceived,回复 SYNACK | +| SynReceived | 相同 SYN | 重新回复 SYNACK,不替换状态或延长截止时间 | +| SynReceived | 匹配的 ACK\|DATA | 完成握手,登记待 accept 一次,停止 SYNACK 定时重传 | +| SynSent | 合法 SYNACK | 确定模式,发送最终 ACK,完成 connect | +| Established 或半关闭 | 重复 SYNACK | 回复最终 ACK,不重新创建流、不改变关闭状态 | +| 已握手、尚未最终释放 | 相同 SYN | 可重发原 SYNACK,不重复 accept、不恢复已关闭方向 | +| 已握手 | 重复最终 ACK | 幂等处理,不重复交接连接 | + +ACK\|DATA 带有效 KCP payload 时仍按 KCP 处理。若它同时完成握手, +不得为了修复握手而重复交付 payload;必须验证 accept 前到达数据时, +现有 KCP 重传和接收路径最终完整交付字节。 + +最终 ACK 丢失的正常恢复: + +```mermaid +sequenceDiagram + participant S as 源端 + participant D as 目标端 + S->>D: SYN + D->>S: SYNACK + S--xD: 最终 ACK 丢失 + Note over S: connect 已返回 + Note over D: 保持 SynReceived + D->>S: 重传 SYNACK + S->>D: 重发最终 ACK + Note over D: 登记待 accept 一次 +``` + +### 5.1 最终 ACK 丢失后立即 FIN + +空流的发送缓冲本来就是空的,connect 返回后可以立即 shutdown;不能 +假设 FIN 一定晚于最终 ACK 到达。 + +可靠 SynReceived 收到匹配模式和连接标识的合法 FIN 时,将它视为 +隐含的握手确认,同时执行: + +1. 完成握手并停止 SYNACK 重传。 +2. 登记待 accept 一次,保存对端发送方向已关闭。 +3. 回复 FIN_ACK;accept 创建流时继承对端关闭状态。 +4. 接收排空后报告 EOF,本端仍可回复数据。 + +不能先进入 Closed 再返回 RST,也不能要求应用先发送非空数据才能 +完成握手。该规则仅用于可靠模式。 + +## 6. FIN、确认与双向关闭 + +### 6.1 本端关闭发送方向 + +保留现有顺序:应用关闭发送队列,发送任务把排队数据交给 KCP,并等待 +`waitsnd == 0`,然后进入 FIN 待确认阶段。重试预算从这一阶段开始计算, +不能从应用开始传输或尚未排空数据时开始计算。 + +FIN 发送入队成功不代表对端已收到。记录待确认状态、下一次重传时间和 +固定截止时间,直到收到有效 FIN_ACK 或连接明确失败。收到对端 FIN +本身不等同于本端 FIN 已被确认。 + +保留现有 AsyncWrite shutdown 的本地关闭语义,不将其返回成功解释为 +对端已确认。关闭确认和最终回收由端点负责;后续超时通过仍存活的流和 +诊断状态表达,不能变成伪造的正常 EOF。 + +### 6.2 收到 FIN 或 FIN_ACK + +- 收到 FIN:记录对端发送方向已关闭,并安排 FIN_ACK。接收任务仍需 + 把已确认接收的数据交给应用,排空后才产生 EOF。 +- FIN_ACK 可以先于应用读完数据发送;它确认端点收到了 FIN,不确认 + 应用已经消费数据。因此接收状态和缓冲生命周期不能随确认提前删除。 +- 重复 FIN:重复应答 FIN_ACK,不重复产生应用 EOF,不重置连接。 +- 收到 FIN_ACK:仅当本端确实已经发送 FIN,且模式、连接标识、flags + 和 payload 全部合法时,停止本端 FIN 重传。 +- 提前、重复或与当前阶段无关的 FIN_ACK 不推进状态、不关闭写方向。 + +发送和接收两个方向分别记录完成情况。现有 Established、LocalClosed、 +PeerClosed、Closed 不足以单独表达“本端 FIN 已发但未确认”,应在现有 +连接状态内补足该信息,不另建一套连接所有权或通用重试框架。 + +### 6.3 正常回收与短暂关闭记录 + +可靠模式最终回收分两步: + +1. 本端 FIN 已确认、对端 FIN 已收到,并且接收任务排空后,释放 KCP + 数据对象、应用通道和数据任务;在原 state_map 保留最小关闭记录。 +2. 固定保留期结束后删除该记录。重复报文不延长保留期。 + +关闭记录仅保存原连接 ID、协议模式、正常关闭原因、回复重复 FIN 所需 +信息和过期时间。它继续回复 FIN_ACK;合法在途 PING/PONG 可以忽略, +不得因现有 Closed 心跳路径而回复 RST。迟到 SYN/SYNACK 不得重新 +accept、创建数据对象或恢复连接;已经完成的数据包可以丢弃。 + +双方同时关闭时,即使只有一方收到 FIN_ACK,先完成的一方也要依靠该 +记录继续回答另一方的 FIN 重传。不能沿用“Closed 且 recv_done 就立即 +删除所有状态”的判断。 + +legacy 继续使用原有回收规则。可靠模式的短暂记录会使状态表在有限 +时间内多保留条目,容量验收必须单独计算,不能继续要求所有内部状态 +在原 15 秒观察点清零。 + +## 7. 定时、队列与失败路径 + +### 7.1 拟定默认参数 + +以下为实现初值,需由定点丢包与延迟测试验证后固定;调整须在验证记录 +说明,不能为了让测试通过而延长业务超时。 + +| 参数 | 初值 / 规则 | +| --- | --- | +| SYN、SYNACK、FIN 首次重传间隔 | 200 ms | +| 后续重传 | 指数退避,单次间隔上限 2 s | +| 源端握手期限 | 调用方现有 `connect(timeout_dur)`,不延长 | +| 可靠目标端握手期限 | 首次收到 SYN 后 60 s,不由心跳或重复包续期 | +| FIN 待确认期限 | 进入 FIN 阶段后 60 s | +| 正常关闭记录保留期 | 65 s,覆盖 FIN 重试窗口并留调度余量 | + +目标端和 FIN 的 60 秒预算是端点控制状态的资源边界,不替代应用更短的 +超时。网络永久不可达时,正确结果是可诊断的超时和最终清理。 + +关闭记录保留期必须不少于另一端允许的最大 FIN 重试窗口加余量,因此 +FIN 期限不能变成各端任意可调且互不知情的参数。若后续版本改变这些 +上界,需要重新审查版本间兼容性。 + +FIN 已确认后立即解除 FIN 重试期限;正常半关闭继续按既有连接存活规则 +处理。仅本端 FIN 未确认才受该期限限制,反向有数据或 PONG 也不能使 +未确认 FIN 无限续期。 + +### 7.2 端点调度与输出队列 + +使用单个端点控制调度任务和现有 state_map;不为每次重传创建任务,不 +扩容或改成无界队列。按最近的控制截止时间唤醒,新状态通过通知更新 +调度;没有待处理控制状态时,避免对全部空闲连接进行高频扫描。 + +状态记录是控制报文待处理事实的来源。输出队列暂满时保留该事实,稍后 +再尝试;不能像原 accept 通知一样丢失后就再无恢复机会。 + +- 用短临界区和非阻塞入队,禁止持有 state_map/conn_map 锁跨 await。 +- 每次实际入队前重新检查状态、连接标识和期限,防止陈旧重传在取消或 + 回收后继续发送。回收不能与入队检查形成复活旧连接的竞态。 +- SYN/SYNACK/FIN 的重传受各自固定预算限制。重复请求所需的 ACK 或 + FIN_ACK 应答可合并为待发状态,不能累计无界报文列表。 +- 应答暂时无法入队时由调度器继续尝试;应答不会获得独立的无限寿命。 + 合并应答记录从首次待发起最多保留 60 秒,重复请求不续期;正常关闭 + 记录中的应答还受该记录更早的截止时间限制。入队成功即撤销待发标记。 + 单纯的应答待发过期只清除该应答,不因此关闭仍合法存活的连接;请求 + 发送方最终依靠自身握手或 FIN 重试期限判定成功或失败。 +- 队列满不能阻塞整个控制调度任务,也不能阻止它处理其他连接的截止 + 时间。不得把“尝试入队成功”统计成“对端确认成功”。 + +### 7.3 取消、错误与连接 ID + +- connect future 取消或超时:撤销该连接的重试与本地所有权;目标端 + 未完成握手的状态按其固定期限退出,重复包不能延长。 +- 有效 RST、端点销毁或 FIN 重试耗尽:停止控制重试、唤醒等待任务, + 保留错误与正常 EOF 的区别。只发送必要的现有重置通知,不建立 RST + 确认和重试协议。 +- 清理失败连接不能重新产生待 accept;迟到的控制响应不能恢复已取消 + 的 connect future。 +- ConnId 分配必须避开所有仍有效的状态,包括正常关闭记录。回绕不能 + 复用仍在保留期内的 ID;会话标识变化后,旧报文不得作用于新会话。 +- 保留期之后无限迟到的包不在可靠恢复窗口内。不得通过永久保存关闭 + 记录来试图覆盖无界的网络延迟。 + +## 8. 兼容性矩阵与降级 + +| 源端 | 目标端 | 预期模式 | 必须满足 | +| --- | --- | --- | --- | +| 旧 | 旧 | legacy | 对照基线,保留既有行为 | +| 新 | 旧 | SYNACK `rsv=0` 确定 legacy | 正常互通,不发 FIN_ACK,不要求目标支持新确认 | +| 旧 | 新 | 从 SYN `rsv=0` 确定 legacy | 正常互通,不主动开启新 SYNACK/FIN 重传语义 | +| 新 | 新 | 本连接握手确认可靠模式 | 握手与关闭在重试预算内恢复 | + +兼容的含义是旧节点仍可互通,且混合组合不因新代码发生行为退化。它不 +意味着旧目标端自动获得 accept 修复,或旧连接获得其未实现的 ACK/FIN +恢复能力。必须把遗留失败与新增失败分别记录。 + +发布或回退影响后续新连接的协商。协议模式固定在连接内,禁止基于后续 +路由公告变化对活跃连接进行升级或降级。旧可执行程序重启后不承担继续 +解释原进程可靠模式连接的义务。 + +## 9. 实施拆分 + +实施使用独立协议修复分支,以 `b37ee660` 及 EasyTier 当前已验证修复 +为基线,不将协议改动混入原 TCP flow-key patch。 + +1. **先建立兼容测试夹具。** 使用实际旧依赖和当前基线端点,确认 `rsv` + 处理、旧 SYNACK、重复 SYN 和旧模式输出门控。明确失败基线。 +2. **实现协商及状态处理。** 增加模式、合法报文校验、可靠模式幂等应答, + 保留 legacy 路径。补上最终 ACK 丢失加立即 FIN 的交叉状态。 +3. **实现统一控制调度及可靠关闭。** 补足握手/FIN 期限、FIN_ACK、正常 + 关闭记录、队列暂满、取消和错误清理。 +4. **完成依赖验证后接入 EasyTier。** 更新依赖 pin 与锁文件,运行真实 + 代理路径、混合版本和跨平台验证,提交可追溯记录。 + +协商声明代表完整可靠控制协议,不能在只实现握手恢复、尚未实现 FIN_ACK +时对外声明 `rsv=1`。中间提交可以用于审阅,但仅完整实现且验证通过的 +依赖版本可被 EasyTier 发布使用。 + +代码主要落点为 kcp-sys 的 `packet_def.rs`、`state.rs`、`endpoint.rs`, +按实际错误传播需要调整 `stream.rs`;不预先拆出通用协议框架。EasyTier +预计只需依赖更新、测试和记录,不需要修改 protobuf 或节点能力公告。 + +## 10. 验证与验收标准 + +### 10.1 确定性协议测试 + +在端点输入/输出之间设置测试用报文过滤器,按连接 ID、报文类型和次数 +精确丢弃、延迟或重复报文。可使用可控时钟测试期限,不依赖随机 netem +才能触发边界,也不通过扩大生产超时使测试通过。 + +| 场景 | 验收断言 | +| --- | --- | +| 分别丢首个 SYN、SYNACK、最终 ACK | 新新组合恢复,同一连接只交接一次 | +| 连续丢多次控制包后恢复链路 | 截止时间内恢复,重试次数与退避符合预期 | +| 最终 ACK 丢失后空流立即 FIN | 目标完成握手、返回 EOF、仍能回复数据,无 RST | +| 最终 ACK 丢失后首个 DATA 到达 | 正确完成握手,应用字节完整且不重复 | +| 重复 SYN/SYNACK/ACK,含半关闭阶段 | 不重复创建连接,不恢复已经关闭的方向 | +| 丢单向 FIN 或 FIN_ACK | 重传恢复;确认不关闭另一方向 | +| 双方同时 FIN,单侧或双侧确认丢失 | 正确排空,正常关闭记录继续应答,最终回收 | +| FIN 先于 accept 或 connect 返回 | 空请求、空响应均正确得到 EOF | +| 接收缓冲尚未被应用读完 | FIN_ACK 不导致缓冲和数据任务提前释放 | +| FIN 已确认后持续半关闭超过 60 s | 反向仍可传输,不受 FIN 期限误杀 | +| 正常关闭记录收到重复 FIN、旧心跳 | 应答或忽略,无 RST,不延长固定保留期 | +| SYN/FIN 永久丢失、控制输出队列长时间满 | 到期退出,无无限任务或状态残留 | +| 队列临时满后恢复 | 待发事实不丢失,其他连接和清理继续运行 | +| connect 取消、RST、端点销毁与重试并发 | 不复活旧状态,等待者正确结束,无伪 EOF | +| ConnId 回绕、会话变化、迟到报文 | 不命中仍在保留期的旧 ID,不污染新会话 | + +### 10.2 真实旧版兼容测试 + +测试夹具必须运行旧依赖的真实 endpoint,不能只在新实现上设置 legacy +标志代替旧版本。首次验收覆盖 `d7427c2`、`b37ee660` 与新实现的组合, +并记录 EasyTier 发布验证实际选用的旧二进制版本和哈希。 + +除四种组合的双向数据与关闭外,额外断言: + +- 旧端收到 SYN `rsv=1` 仍按旧协议回复 SYNACK `rsv=0`。 +- 重复 SYN 不引出混合连接新增 RST;legacy 输出门控保持原有行为。 +- 混合组合捕获不到 FIN_ACK;旧端不必识别任何新控制语义。 +- PING/PONG、错误连接 ID、非法 flags 或未知 SYNACK 模式不能确认能力。 +- 新连接在节点升级、回退后重新选择正确模式,不沿用节点级缓存。 +- 旧版原有半关闭或丢包失败作为对照保留,不将它们写成新协议已通过。 + +### 10.3 EasyTier 流量与平台验证 + +Linux 需要 root 的测试在现有 `rust` 容器中运行。复用正常 target,权限 +问题通过修复所有权解决,不另开编译目录。 + +- 完整三节点组合及 ACL、配置更新、端口转发、断连测试。 +- TCP/KCP/QUIC × 内核/smoltcp 六种模式,检查未修改协议没有退化。 +- 每组合 300 次短连接、16 并发、5 秒 socket 超时;KCP 每栈追加三轮, + 与已有每版本 2,400 次 KCP 结果对照。 +- 空请求、256 KiB 请求后 EOF、1 MiB 响应,以及服务端先半关闭后客户端 + 才发送 256 KiB 的反向场景,均逐字节核对。 +- `netem delay 10ms 3ms loss 1%` 下持续双连接与短连接;定点控制包丢失 + 由协议测试证明,随机丢包实验用于验证整体行为。 +- 新旧两方向真实二进制互通;无丢包与丢包结果分别记录。 +- Linux、macOS、Windows 原生依赖测试及网关测试;格式与严格 Clippy。 + 缺少组件、未运行的项目不得记录为通过。 + +### 10.4 未解释的反向超时 + +复现脚本必须记录超时时已收到的字节数、是否收到 EOF、源端口、连接 ID, +并在双方记录 FIN/FIN_ACK 的入队、发送和接收时间。补足原失败缺少的 +证据,区分数据缺失、关闭通知缺失、队列延迟和状态处理错误。 + +恢复实际链路的 FIN 丢失,只能证明这一类注入故障已修复;仍需解释原 +异常或明确保留未定位项。不得用重跑通过覆盖原始失败。 + +### 10.5 资源与交付门槛 + +分别统计:应用代理连接、KCP 数据对象、半开握手、未确认 FIN、正常关闭 +记录、FD 和 RSS。保留期内的轻量记录属于设计成本,不能混同于活跃连接 +泄漏;其数量约受每秒关闭连接数乘以保留期约束,需实测内存成本。 + +静止并超过握手/FIN 期限及关闭记录保留期后,所有应回收状态必须消失。 +使用重复负载周期检查资源是否持续累积,并检查端点空闲及大量并发时 +控制调度的 CPU 成本。不能仅凭一次 FD 回到基线宣布不存在泄漏。 + +交付必须满足:确定性恢复用例通过;实际旧版兼容用例没有新增失败; +错误路径有界清理;原有矩阵没有新退化;所有异常如实保留。完整最终 +代码 diff 由独立子代理审查,只处理高置信度真实缺陷。验证记录绑定 +提交、依赖 revision、二进制哈希、准确命令及原始日志。 + +## 11. 当前进度与证据 + +已完成:连接级版本协商、握手与 FIN 恢复、FIN_ACK、重复控制处理、 +正常关闭记录、队列饱和处理与取消清理。kcp-sys 最终提交为 +`3ef5c4161faf99940f3ed51efd43cef0cbc02b4f`,EasyTier 最终依赖接入为 +`ee02b8f7`。首次验收已测试真实旧实现 `d7427c2` 和 `b37ee660`。 +后续测试整理仅保留 `d7427c2` 长期基线;当前维护版本与依赖 pin 见 +[实现与验证记录](kcp-control-reliability-validation-2026-09-14.md)。 + +三平台原生依赖测试、Linux 完整矩阵及实际流量的本轮结果单独记录, +不沿用上一轮局部修复的通过数。最终代码审查没有 blocker / major, +一项首次 SYNACK 调度竞态 minor 按用户规则记录待办。 + +仍需保留的边界:原反向超时没有足够证据做最终归因;混合版本保留 +legacy 的控制恢复限制;有限负载和应用级资源采样不能代替生产规模 +长期内存、CPU 和容量验证。所有实测异常及未执行项见关联验证记录。 + +此前调查原始产物位于: + +```text +/data/project/proxy-close-validation-20260914/ +``` + +关键证据:`remaining-loss-handshake.md`、`loss-handshake-evidence.log`、 +`delivery-reverse-timeout.md`、`fin-send-path-old-new.txt`、 +`delivery-manifest.json`、`traffic-summary.md`。上述绝对路径是本机 +调查产物位置;仓库读者可通过关联验证记录了解结论与证据限制。 diff --git a/docs/kcp-control-reliability-validation-2026-09-14.md b/docs/kcp-control-reliability-validation-2026-09-14.md new file mode 100644 index 00000000..f20173fa --- /dev/null +++ b/docs/kcp-control-reliability-validation-2026-09-14.md @@ -0,0 +1,296 @@ +# KCP 控制报文可靠性实现与验证(2026-09-14) + +本文记录 +[协议设计](kcp-control-reliability-design-2026-09-14.md) +的实现与验证,独立于此前的 TCP flow-key 和半关闭局部修复。 + +## 兼容测试维护整理 + +当前维护版本为 kcp-sys `268533568d734ae89dc89603078da3ca522effe1`。 +保留 `d7427c22` 作为长期兼容基线,移除中间版本 `b37ee660` 的 +`kcp-sys-baseline` dev-dependency 与三项重复用例;单次使用的宏展开为 +普通测试函数。四项旧版兼容测试继续覆盖能力协商、双向数据、半关闭及 +重复 SYN;15 项 library 与 12 项协议回归不变。 + +整理后 Linux 共 31/31 测试通过,格式与严格 Clippy 通过。EasyTier +同步依赖 pin 与锁文件,并通过 `cargo +1.95 check --locked -p easytier +--features full`。此次只整理测试与测试依赖,协议源码未改变。 +历史三平台 34/34 和全部流量结果仍属于下述 `3ef5c416` 实现验证, +没有重写为整理后版本的执行结果。中间基线作为调查证据保留在本文。 +本次日志位于 `/data/project/kcp-compat-cleanup-validation/`。 + +## 版本与实现 + +- kcp-sys 基线:`b37ee660fb70bb6d816fb8bbc08b140e55e7218b`。 +- kcp-sys 初版实现:`c84733d4479b40a299d51f4c5b8bb02ccacadc68`。 +- kcp-sys 最终实现:`3ef5c4161faf99940f3ed51efd43cef0cbc02b4f`,补齐 + 可靠模式下未知连接 RST 输出队列饱和时的非阻塞处理。 +- EasyTier 基线:`851e7523`;接入提交:`4fedbdd1`、`ee02b8f7`, + 仅修改依赖 pin 和 Cargo.lock。 +- 最终真实流量二进制:`final-easytier-core`,SHA-256: + `962f2480cb611ceb6cab293560d0dfb337218594cdc298a08f0e83801a2d17b8`。 + 它在接入提交前构建,源码及依赖内容与该提交一致;识别产物以哈希和 + 依赖 revision 为准,不单凭内嵌的 EasyTier git 版本字符串。 + +依赖已发布到 `EasyTier/kcp-sys` 的独立分支 +`fix/control-reliability-20260914`,远端引用核对为上述最终 revision。 +EasyTier 任务分支为 `fix/kcp-control-reliability`。 + +实现保留 14 字节 header,通过 SYN/SYNACK 的 `rsv` 协商可靠控制模式。 +收到旧 SYNACK 后固定使用 legacy,旧源端连接新目标端也从头使用 legacy。 +只有协商为可靠模式的连接使用独立 FIN_ACK 和新增恢复逻辑。 + +可靠模式在原连接 ID 上恢复 SYN、SYNACK 和 FIN;重复握手不重复交接。 +最终 ACK 丢失后的空 FIN 可以同时确认握手并报告对端半关闭,反向仍能 +发送数据。FIN_ACK 只确认收到 FIN,不关闭本端发送方向,也不提前释放 +尚未被应用读取的接收数据。 + +控制请求和待发送应答在原连接状态中保存,由端点统一调度。输出队列 +暂满不会丢失待发事实或延长固定期限;取消 connect 同步撤销本地状态。 +状态检查、控制入队和清理保持一致锁序,不持锁跨 await。 + +初始重传间隔 200 ms、指数退避至 2 s;源端使用调用方连接超时,目标端 +握手及 FIN 待确认期限为 60 s。正常双向关闭且接收排空后释放数据对象, +原状态表保留 65 s 的关闭记录,继续应答重复 FIN,并避免迟到心跳引出 +RST。FIN 一旦确认,正常半关闭不受 FIN 重试期限限制。 + +## 自动化结果 + +| 验证 | 结果 | +| --- | --- | +| Linux kcp-sys | 15 library + 12 协议回归 + 7 真实旧依赖兼容,34/34 | +| macOS kcp-sys | 相同 34/34,两个 example target 通过 | +| Windows kcp-sys | 相同 34/34,两个 example target 通过 | +| Linux 两个 example target、格式、严格 Clippy | 通过 | +| macOS 格式、严格 Clippy | 通过 | +| Windows 格式、Clippy | 所选 stable 缺少组件,未执行 | +| Linux 原生网关测试 | 200/200 | +| macOS、Windows 原生网关测试 | 各 200/200 | +| EasyTier 完整三节点及补充集成测试 | 276/276,834.639 s | +| EasyTier Linux 严格 Clippy | 通过 | + +macOS 首次获取依赖遇到 GitHub TLS 错误,随后导入本机真实 Git 对象与 +checkout,在依赖 revision 不变的情况下离线测试;没有改成 path 依赖, +也没有用新实现替换旧依赖。Windows 原始日志为 UTF-16LE,另保存 UTF-8 +副本。两平台复用原 target,没有另开编译目录绕过权限问题。 + +协议修改集中在 kcp-sys,EasyTier core 源码没有变化。原生网关验证检查 +与现有相同 core 源码的兼容行为,不能替代上面的真实新旧 KCP 端点测试。 +远端应用 manifest 仍锁定旧依赖 `d7427c2`,网关命令只选择 easytier-core; +关键转发文件 `tcp_proxy_service.rs` 哈希与本地一致。因此不能将这组 +网关测试描述为完整原生应用已经接入最终新依赖。 + +## 失败到通过的对照 + +- 丢弃第一份最终握手 ACK:基线客户端 connect 返回,服务端 accept + 超过 5 s 仍未完成;实现后约 0.25 s 恢复,并成功发送服务端 greeting。 +- 新源端连接真实旧依赖:基线仅因 SYN 尚未声明协商能力而未达到新协议 + 测试要求;实现后 SYN 提出 `rsv=1`、旧 SYNACK 回复 `rsv=0`,后续 + 数据与关闭均走 legacy。该项验证协商功能,不把它描述为旧版互通 bug。 +- 原有 RST 单测曾在新连接上注入 `rsv=0` 的合成 RST,因可靠模式拒绝 + 不匹配模式而超时;测试改为注入该连接实际协商版本的 RST 后通过。 + 模式不匹配的报文不能用于证明正常 RST 错误传播失败。 +- 新增双向大缓冲夹具最初单次写入超过既有 KCP send 的分片限制,出现 + `Err(-2)`。调整为与原有测试一致的 16 KiB 分块,仍验证双向各 + 200 KiB 总数据和原 5 s 期限;本补丁没有修改既有单次大写入限制。 + +- 最终补查发现初版 `c84733d` 对未知可靠 FIN 的 RST 仍使用阻塞发送。 + 填满输出队列后,新 SYN 不能在 100 ms 内进入状态表;改为可靠模式 + 的无状态应答使用 try_send 后通过。需要重试的有状态控制仍由状态表 + 保管;legacy 发送路径不变。 + +原始失败日志保留,未使用扩大业务超时或重跑通过覆盖失败记录。 + +## 协议与生命周期覆盖 + +12 项公有 API 协议回归使用真实端点及按报文类型过滤的链路: + +- 连续丢 SYN、连续丢 SYNACK、丢最终 ACK 后的服务端 greeting。 +- 所有空最终 ACK 均丢失,空 FIN 直接完成握手,EOF 后仍可回复。 +- 单向 FIN 或 FIN_ACK 丢失,另一方向仍可传输。 +- 双向 200 KiB 缓冲、双方首个关闭确认丢失、延后读取与排空。 +- 半关闭后重复 SYN/SYNACK/ACK;握手 ACK 丢失时 DATA 乱序、重复。 +- FIN 确认后推进 61 s,再恢复实际时钟,仍可反向传输。 +- PONG 保留 `rsv=1`、非法 SYNACK flags、非空 SYNACK 均不能确认模式。 +- 未知 SYNACK version 返回 `InvalidProtocolVersion`,不启用新模式。 + +15 项 library 测试包含原有 10 项以及五项生命周期与队列测试:取消 connect +立即释放状态;重复 SYN 与满输出队列不延长半开期限;未确认 FIN 超时 +唤醒读端并释放状态;正常关闭记录应答迟到包、不引出 RST、不被延长, +保留期间跳过相同 ConnId,期满后删除旧记录而不影响新连接;满输出 +队列下未知可靠 FIN 不能阻止后续 SYN 进入握手。 + +可控时钟只用于测试期限,不修改生产时间常量。数据测试仍通过实际 KCP +发送、接收及 AsyncRead/AsyncWrite 路径。 + +## 真实旧依赖兼容 + +dev-dependency 固定并实际运行 `d7427c2`、`b37ee660` 两个历史实现, +没有用新代码上的 legacy 开关模拟旧端。 + +7 项测试覆盖旧 SYN 处理探针和两个基线的新旧双向连接、256 KiB 双向 +数据、半关闭、重复 SYN。报文捕获断言:新源初始 SYN 可为 `rsv=1`, +其余混合连接报文均为 `rsv=0`;没有 FIN_ACK、新增 RST 或重复 accept。 + +`d7427c2` 原有正常关闭被报告为 BrokenPipe 的行为仍在对照中保留。 +兼容意味着旧节点仍可互通,不意味着它自动获得新协议的恢复能力。 + +## 实际代理流量 + +实验使用独立 namespace,底层 UDP, +覆盖 TCP/KCP/QUIC 与内核/smoltcp 六种模式。固定每轮 300 次短连接、 +16 并发、5 s socket 超时,逐字节校验大小请求、纯空请求及反向半关闭。 + +验证分为两个阶段;不将初版流量统计冒充最终版本结果。 + +### 初版 c84733d + +二进制 `protocol-easytier-core` 的 SHA-256 为 +`f41ff98762ed1ce8691d8f83b8c9e47acc4548c76432a17cd2daba36e2cc7b80`。 + +- 六模式主矩阵:1,800 次短连接、12 次大小请求半关闭、6 次反向 + 半关闭、192 次纯空请求,全部通过。 +- KCP 每栈追加三轮,合并主矩阵共 2,400 次 KCP 短连接,全部通过。 +- `netem delay 10ms 3ms loss 1%` 六模式:1,800 次短连接、12 次大小 + 半关闭、6 次反向、192 次纯空及每组合 15 s 双连接持续流量,全部通过。 +- 新旧双方向各六模式:3,600 次短连接、24 次大小半关闭、384 次纯空 + 和每组合 3 s 持续流量通过;**反向半关闭为 11/12,存在一次失败**。 + +旧端为依赖 b37 的 `delivery-easytier-core`,SHA-256 为 +`4635810ac9a536258f9b7ff606e5e1a5eb0243cb0f187825ba6f721d92dfc4fe`。 +失败发生在新 source → 旧 destination 的 KCP/smoltcp,连接 +`conv=3560055726`、源端口 `40748`:客户端已经收到完整 1 MiB,随后 +等待 EOF 超过 5 s;因此未进入后发 256 KiB 阶段。旧目标端在 +10:11:08.953 进入 LocalClosed,直到客户端超时才见客户端方向关闭。 +DEBUG 不能确定旧端 FIN 后续的发送、到达或处理点,不能直接定性为 +FIN 丢失,也不能把混合版本测试写成全部通过。 + +另一次有界 KCP/kernel 丢包 TRACE 测试中,300 次短连接及全部半关闭 +通过。实际捕获的控制事件全部使用 version 1:连接 `1468864199` 的 +重复 SYNACK 相隔约 200 ms,目标随后收到 ACK,业务完成;连接 +`1468864161` 的重复 FIN 后收到 FIN_ACK。它们证明实际链路使用了 +新控制恢复路径;单轮成功不代表任意丢包条件下都能成功。 + +### 最终 3ef5c416 + +- 六模式主矩阵全部通过:1,800 次短连接、12 次大小半关闭、6 次反向 + 半关闭、192 次纯空请求。 +- 六模式随机丢包全部通过:相同流量规模,另每组合 15 s 双连接持续 + 逐字节回显。原 5 s socket 超时不变,没有添加应用重试。 +- 混合版本无丢包,两连接方向 × 两个 KCP 栈全部通过:1,200 次短连接、 + 8 次大小半关闭、4 次反向及 128 次纯空请求。该轮通过不覆盖初版 + 混合实验中的 EOF 超时记录。 +- KCP 每栈追加三轮全部通过;合并主矩阵为 **2,400/2,400** 次 KCP + 短连接;追加轮次的 12 次大小、6 次反向和 192 次纯空也全部通过。 +混合丢包仍走 legacy,下表为单轮**失败数**,不重跑取最好结果。 +每组合为 300 短连接、2 次大小半关闭、1 次反向、32 次纯空和 3 s 持续流量。 + +| 连接方向/栈 | 短连接失败 | 大小失败 | 反向失败 | 纯空失败 | 持续失败 | +| --- | --- | --- | --- | --- | --- | +| b37 → 最终新 / kernel | 6 | 0 | 0 | 16 | 0 | +| b37 → 最终新 / smoltcp | 2 | 0 | 0 | 13 | 0 | +| 最终新 → b37 / kernel | 7 | 1 | 0 | 16 | 0 | +| 最终新 → b37 / smoltcp | 1 | 0 | 0 | 14 | 0 | + +唯一大小用例失败停在 greeting 阶段,5 s 超时;纯空失败为 response +阶段收到空 EOF。为判断兼容回退,另跑一次相同条件、相同业务参数的 +b37 → b37 旧旧对照:kernel 短连接失败 2/300、纯空失败 12/32; +smoltcp 分别为 0/300、16/32。两栈的大小、反向与持续项目均通过。 + +旧旧 TRACE 确认了一条失败链:kernel 的 `conv=3046346801`、空请求 +`id=3`、源端口 `35958`,目标在 10:30:44.112 先收到 FIN,直接 +Closed 并发送 RST;10:30:44.115 才收到最终 ACK|DATA。源在 .113 +收到 RST,客户端约 37 ms 后得到零字节 EOF。`src/state.rs` 的 +SynReceived + FIN → Closed/RST 路径本次没有修改;可靠模式单独支持 +FIN 完成握手,legacy 按已批准设计保留原行为。 + +这条时序证明上述失效路径在旧旧连接也存在,不能由随机样本的失败率 +断言所有混合版本失败均已归因或已经排除一切回退。混合实验原 DEBUG +不足以逐包归因全部失败。兼容保证旧节点可按原协议互通;控制恢复能力 +需要两端都协商为 version 1,混合部署仍不能获得这一保证。 + +最终丢包轮次直接启用 KCP TRACE:kernel 捕获 5 个成功连接收到重复 +SYNACK、4 个收到重复 FIN;smoltcp 分别为 2 个、5 个。捕获的控制 +报文均使用 version 1。详见 `final-control-recovery-*-summary.json` +及对应 `evidence.log`;这份证据直接绑定最终二进制。 + +资源观察区分应用代理条目与内部 KCP 状态。CLI 不直接暴露依赖内部 +关闭记录,不能用 CLI 条目清零证明内部记录已经删除;内部期限由 +生命周期测试验证。实际流量保存起始、结束、15 s、80 s 的 FD/RSS 与 +应用代理条目。65 s 关闭记录属于设计成本,不能沿用旧的 15 s 内部 +状态必须清零的断言。 + +最终每栈三轮连续负载的资源结果如下。a 为 source,b 为 destination。 +该追加轮次的 FD 全程未增长;应用 proxy 条目在 15 s 与 80 s 均为零。 + +| 栈/节点 | FD 前后 | RSS 前 → 80 s(KiB) | 负载 CPU 秒 / 墙钟秒 | 静置 15–80 s 单核 CPU | +| --- | --- | --- | --- | --- | +| kernel/a | 15 → 15 | 46,944 → 50,792 | 2.20 / 2.697 | 0.29% | +| kernel/b | 15 → 15 | 46,580 → 49,344 | 1.12 / 2.694 | 0.31% | +| smoltcp/a | 14 → 14 | 46,784 → 50,352 | 4.39 / 5.378 | 0.49% | +| smoltcp/b | 14 → 14 | 46,024 → 49,060 | 2.32 / 5.368 | 0.48% | + +CPU 来自 `/proc/PID/stat` 的 user/system ticks 与单调时钟差,包含 +路由、心跳、代理、日志等全部进程工作。产物为 debug 构建,表中负载 +CPU 也不是纯协议开销或吞吐基准。RSS 未回到起点,无法由这些采样区分 +分配器缓存与其他长期对象,也不能推算每条关闭记录的精确内存成本。 +内部对象期限由单测验证,生产规模的逐对象内存与长期容量验证仍未完成。 +主矩阵 kernel/b 在 15 s 采样曾由 15 个 FD 暂升为 16,80 s 恢复为 15; +该瞬时变化同样保留在原始报告,未作为持续增长处理。 +原始记录为 `final-kcp-repeat-kcp-*-resources.json`,换算另存 +`final-resource-summary.json`。 + +## 审查、复现与限制 + +初轮独立子代理只读审查 `b37ee660..c84733d`,没有高置信度缺陷发现。 +最终由新子代理审查完整 `b37ee660..3ef5c416` 及 EasyTier 最终 pin, +未发现 blocker / major;重点检查旧版输出规则、握手与关闭交叉状态、 +队列、取消清理和锁序。 + +最终审查记录一项 **minor / high confidence** 待办: +`kcp-sys/src/endpoint.rs:922` 新 SYN 路径先 notify 再插入连接状态, +多线程时可能先消费通知并漏过首次 SYNACK 调度。正常源端约 200 ms 后 +重发 SYN 即可恢复;若后续 SYN 未到达,则等待约 10 s 周期扫描。状态 +不会丢失,也不会永久阻塞,但特别短的 connect 期限可能超时。按用户 +minor 默认记录的规则保留;后续最小改动是将通知放到状态插入之后。 + +主要命令: + +```sh +# kcp-sys 工作树 +cargo test --all-targets +cargo fmt --all -- --check +cargo clippy --all-targets -- -D warnings + +# EasyTier,Linux root 集成测试在 rust 容器内执行 +cargo +1.95 test --locked -p easytier-core \ + --features proxy-smoltcp-stack --lib gateway:: +cargo +1.95 nextest run --locked -p easytier --features full --lib \ + -E 'test(subnet_proxy_three_node_test) | test(subnet_proxy_half_close_test) | test(acl_rule_test_inbound) | test(acl_rule_test_subnet_proxy) | test(proxy_three_node_disconnect_test) | test(config_patch_test) | test(port_forward_with_inbound_default_drop_acl_test)' \ + --test-threads 1 --no-fail-fast +cargo +1.95 clippy --locked -p easytier --features full \ + --lib --tests -- -D warnings +``` + +本机原始日志、脚本、JSON 与二进制: + +```text +/data/project/kcp-control-validation-20260914/ +``` + +主要证据为 `final-ack-red.log`、`final-ack-green.log`、 +`final-dependency-tests.log`、`final-dependency-clippy.log`、 +`unknown-close-full-output-red.log`、`legacy-final.log`、 +`negotiation-lifetime-regressions.log`、`final-native-*-tests.log`、 +`linux-gateway.log`、`integration-final.log`、`app-final-clippy.log`。 +`integration-protocol.log` 是初版的中断轮次,切换最终实现后重新执行完整 +矩阵,不计入最终通过数。 +流量的准确 argv、产物 hash 和逐连接结果另存该目录。完整流量汇总为 +`traffic-summary.md`,混合丢包对照为 `mixed-loss-analysis.md`, +旧旧逐包证据为 `legacy-fin-before-ack-evidence.log` 与 +`legacy-fin-before-ack-client.json`。测试环境清理记录为 `cleanup.json`。 + +此前单次反向超时的具体丢包点仍不能由旧 DEBUG 日志倒推出。定点 FIN +丢失测试证明该类失效现在能够恢复,不能因此改写原事故的根因结论。 +Windows/macOS 原生 TUN 端到端、生产规模长期负载、吞吐与容量上限仍 +不是这些有限测试能够证明的事项。 diff --git a/docs/tcp-proxy-flow-key-validation-2026-09-13.md b/docs/tcp-proxy-flow-key-validation-2026-09-13.md new file mode 100644 index 00000000..f269f176 --- /dev/null +++ b/docs/tcp-proxy-flow-key-validation-2026-09-13.md @@ -0,0 +1,137 @@ +# TCP proxy flow-key 验证记录(2026-09-13) + +本次验证没有发现只在修复版本出现的行为退化。同源端口、不同目标的 +并发连接在六种代理模式下均由父提交的超时变为成功。验证中仍有半关闭 +失败和 KCP 突发短连接超时,父提交也存在这些现象,不能将结果描述为 +所有场景均无异常。 + +## 版本与范围 + +- 修复版本:`eb83655958be932d6de34e090dc361b6f3ba3393`。 +- 父提交:`e0bdb516b6dc8a654940efbe12960dbfa846f424`。 +- 本次提交仅新增测试与记录,生产代码保持上述修复版本的内容。 +- Linux 测试在现有 `rust` 容器内运行,Rust 1.93.1。 +- macOS arm64、Windows x64 原生测试使用 Rust 1.95.0。 + +父提交使用独立 worktree,但复用当前 worktree 的 `target`。第一次父提交 +网关测试意外复用了修复版本的构建缓存,结果已排除。清理相应 package +的构建产物后重新编译,确认父提交网关测试为 189 项,修复版本原有 +192 项。真实流量实验还通过 CLI 核对两端实际运行的版本号。 + +## 自动化验证 + +| 验证 | 结果 | +| --- | --- | +| 修复版本完整 `subnet_proxy_three_node_test` 矩阵 | 256/256,714.800 秒 | +| ACL、端口转发 ACL、配置更新、代理断连 | 父提交与修复版本各 14/14 | +| 父提交 Linux 网关测试 | 189/189 | +| 新增测试后的 Linux 网关测试 | 198/198 | +| 新增测试后的 macOS 原生网关测试 | 198/198 | +| 新增测试后的 Windows 原生网关测试 | 198/198 | + +完整三节点矩阵覆盖 TUN/no-TUN、普通/公共中继、源端 KCP/QUIC 开关、 +目标端 KCP/QUIC 开关及对应输入禁用组合。每个组合检查映射子网地址、 +真实子网地址和节点虚拟地址上的 ICMP、TCP、UDP。 + +命令(Linux 需在容器内运行): + +```sh +cargo nextest run -p easytier --features full --lib \ + subnet_proxy_three_node_test --test-threads 1 --no-fail-fast + +cargo test -p easytier-core --features proxy-smoltcp-stack --lib gateway:: + +cargo fmt --all -- --check + +cargo clippy -p easytier-core \ + --features proxy-smoltcp-stack,ring-crypto --lib --tests -- -D warnings +``` + +格式检查和上述严格 Clippy 检查均通过。仅启用 `proxy-smoltcp-stack` +时,严格 Clippy 被未修改的 `tunnel/encrypt/mod.rs` 中 +`assert_interoperable` 未使用警告阻断;增加 `ring-crypto` 会启用调用 +该函数的现有后端互操作测试。本次没有修改或屏蔽该警告。 + +14 项补充集成测试来自以下测试组,使用各版本独立保存的测试二进制, +逐项 `--exact` 执行,避免不同进程同时操作相同的测试 network namespace: + +- `acl_rule_test_inbound`:4 项。 +- `acl_rule_test_subnet_proxy`:4 项。 +- `port_forward_with_inbound_default_drop_acl_test`:3 项。 +- `config_patch_test`:1 项。 +- `proxy_three_node_disconnect_test`:2 项。 + +三节点矩阵与流量实验在新增测试前完成;新增测试不改变生产代码,随后 +在三平台重新执行完整网关测试。macOS、Windows 的结果不包含原生 TUN +端到端验证。 + +## 新增回归测试 + +测试位于 `easytier-core/src/gateway/proxy/tcp_proxy_engine.rs`: + +1. SYN 在 accept 前后重传,均保留转换端口,且不会重复 accept。 +2. 旧连接处于 ClosingSrc、ClosingDst、Closed 时,被同一流的新连接 + 替换;旧连接清理不会删除新映射,反向地址、端口和校验和仍正确。 +3. 过期 SYN 同时释放两个索引,已 accept 的连接不受 SYN 超时清理影响。 +4. 转换端口计数器回绕后跳过零、已占用端口和监听端口。 +5. 实际填满同一源 IP 的 65,534 个转换端口后,新 SYN 被丢弃;其他源 + IP 仍可建立映射,释放一个条目后该端口可复用,clear 后也可重新分配。 +6. 八个线程同时处理同一流,仅建立一个映射并成功 accept 一次。 + +端口池测试验证实际容量和恢复行为,没有将 debug 构建的耗时用作生产 +性能阈值,也没有通过延长超时或添加重试改变被测行为。 + +## 真实流量对照 + +使用独立的 `flow_val_a`、`flow_val_b` namespace,物理链路为一对 veth。 +虚拟地址为 `10.251.92.1/24` 和 `10.251.92.2/24`;目标端将 +`192.0.2.0/24` 映射到 `198.18.0.0/24`。两个 TCP 服务监听真实地址 +`.10`、`.11` 的 23456 端口,返回服务地址标记,并逐字节回显、校验数据。 +底层隧道使用 UDP,两个节点均使用四线程运行时。 + +每个版本分别验证普通 TCP、KCP、QUIC 与内核/smoltcp 的六种组合。 + +| 场景 | 父提交 | 修复版本 | +| --- | --- | --- | +| 源 IP/端口相同、两个不同目标的并发连接,各进行 30 轮回显 | 6/6 超时 | 6/6 通过 | +| 同一四元组 RST 后重连,30 次,间隔 20ms | 6/6 通过 | 6/6 通过 | +| 普通 TCP、QUIC 的短连接,各 300 次、16 并发 | 各组合 300/300 | 各组合 300/300 | +| KCP 内核模式短连接,300 次、16 并发 | 4 次超时 | 4 次超时 | +| KCP smoltcp 模式短连接,300 次、16 并发 | 12 次超时 | 11 次超时 | +| 关闭客户端写端后等待服务端响应 | 6/6 无响应数据 | 6/6 无响应数据 | + +各组合完成流量后等待 15 秒。修复版本两端的代理条目均回收为零,FD +数量回到接近起始水平(采样期间 RPC 连接带来约一个 FD 的差异)。这只能 +证明本次有限负载下的回收行为,不能证明长期内存占用没有增长。 + +KCP 的超时阈值为 5 秒,尚未定位这些突发短连接超时的根因;两版都有 +失败不等于已证明每次失败属于同一根因。半关闭实验的服务端在 EOF 后 +返回响应;未修改的 `copy_bidirectional_no_shutdown` 在任一方向结束 +后即退出转发,与两版都丢失该响应的现象一致。本次不修复这两类问题。 + +## 丢包与混合版本 + +- 在客户端 veth 出方向施加 `netem delay 10ms 3ms loss 1%`。 +- 两版各六种模式,每个组合保持两个连接持续回显至少 15 秒,所有数据 + 均通过逐字节校验;共 12/12 通过。 +- 无 netem 时,旧源端/新目标端、新源端/旧目标端,各覆盖六种模式。 + 两个连接持续回显至少 2 秒,共 12/12 通过。 +- 混合版本验证使用不同客户端源端口,验证正常互通;它不意味着仍运行 + 旧代理引擎的一端也获得了同源端口冲突修复。 + +## 证据与限制 + +本机原始日志、流量脚本、各阶段 JSON、二进制与 SHA-256 清单保存在: + +```text +/data/project/tcp-flow-validation-20260913/ +``` + +主要记录为 `manifest.json`、`current-matrix-summary.log`、 +`*-extra.json`、`traffic-matrix.json`、`loss-interop.json` 和 +`current-final*gateway.log`。矩阵记录是最终摘要,不是完整逐项日志。 +实验进程、独立 namespace 和 netem 均已清理。 + +未覆盖 Windows/macOS 原生 TUN 路径、生产防火墙/conntrack 规则兼容性、 +长期运行、吞吐回归基准,以及接近容量极限时的真实内核连接负载。 +端口耗尽已在引擎级验证,但不能替代上述生产容量测试。 diff --git a/docs/tcp-proxy-half-close-validation-2026-09-14.md b/docs/tcp-proxy-half-close-validation-2026-09-14.md new file mode 100644 index 00000000..c4845da2 --- /dev/null +++ b/docs/tcp-proxy-half-close-validation-2026-09-14.md @@ -0,0 +1,178 @@ +# TCP 代理半关闭与 KCP 短连接修复验证(2026-09-14) + +本次修复上一轮验证发现的半关闭失败、无丢包 KCP 突发连接交接丢失, +以及跨平台验证暴露的握手前心跳竞态。这些是转发与连接生命周期的局部 +逻辑问题;没有改变 KCP 报文格式,也没有增加重试或放宽业务超时。 +丢包网络下另有握手最终 ACK 丢失的问题,本次未修复,见下文。 + +## 版本与根因 + +- EasyTier 基线:`24572b49`;生产修复:`03fa375e`、`fd6e2a3e`。 +- kcp-sys 基线:`d7427c22d764deb1860a7d37acc446ed5033464c`。 +- 最终依赖:`b37ee660fb70bb6d816fb8bbc08b140e55e7218b`, + `Cargo.toml` 和 `Cargo.lock` 均锁定此 Git revision。 + 已发布到 kcp-sys 的 `fix/accept-half-close-20260914` 专用分支, + 便于其他机器获取锁定的依赖;未合并主分支。 +- 旧版真实流量对照:`eb836559`,已包含 flow-key 修复,尚无本次修复。 +- 最终流量二进制:`delivery-easytier-core`,SHA-256: + `4635810ac9a536258f9b7ff606e5e1a5eb0243cb0f187825ba6f721d92dfc4fe`。 + 构建内容与最终生产代码一致;构建发生在依赖 pin 提交前,内嵌版本号 + 不用于区分本次产物,以 SHA-256 和依赖 revision 为准。 + +修复内容: + +1. 原转发函数在任一方向 EOF 后退出,取消另一方向,服务端在请求 EOF + 后返回的响应因此丢失。改用 Tokio 双向拷贝传播写端关闭,继续转发 + 反向数据,直到双向完成或发生错误。 +2. KCP 容量为 4 的 accept 通知队列满时,已建立连接失去交接机会。 + 在连接状态中保留待 accept 标记,队列为空时领取待交接连接;保留 + 原有有界通知队列和串行领取,防止重复交接。 +3. KCP 将正常关闭当作读错误,并可能在 FIN 后丢弃已经确认接收、尚未 + 交给应用的数据。现在先排空接收数据再返回 EOF,RST 和端点销毁仍 + 返回错误;双向关闭清理等待接收任务完成。EOF 判断和过期清理的 + 状态读取顺序也一并修正,避免并发数据到达或心跳更新被旧判断覆盖。 +4. FIN 可以先于 accept 或 connect 返回到达。创建流时继承该关闭状态, + 使纯空请求、纯空响应也能得到 EOF。半关闭连接继续响应心跳。 +5. Windows 原生测试捕获到 PING 先于 SYN 到达,对端因连接未知返回 + RST;旧依赖也复现相同时序。周期心跳现在只覆盖已建立和半关闭状态。 + +## 回归覆盖与自动化结果 + +| 验证 | 结果 | +| --- | --- | +| Linux、macOS arm64、Windows x64 原生网关测试 | 各 200/200 | +| 最终 kcp-sys 三平台原生测试 | 各 10/10,两个 example target 通过 | +| 最终完整三节点矩阵及补充集成测试 | 276/276,848.628 秒 | +| EasyTier 格式检查、Linux 严格 Clippy | 通过 | +| kcp-sys Linux/macOS 格式检查、严格 Clippy | 通过 | + +网关新增两个测试,分别从两端先关闭写方向;使用空请求、32 KiB 请求 +及 64 KiB 响应,接收方等到 EOF 才返回响应。旧转发函数两项均在 5 秒 +超时,修复后通过。原生网关测试不依赖应用层 KCP,心跳修复没有改变 +这些测试的生产代码。 + +新增三节点半关闭测试覆盖 TCP/KCP/QUIC × 内核/smoltcp 六种模式, +各检查两个关闭方向、空请求和 64 KiB 请求、128 KiB 响应。完整集成 +集合共 276 项:原有三节点矩阵 256 项、新增 6 项、ACL/配置更新/ +端口转发/断连 14 项,均在已有 `rust` 容器中串行运行。 + +依赖回归覆盖:32 次握手全部完成后才 accept;200 KiB 缓冲数据后的 +EOF/响应;双向关闭后跨越清理周期再读数据;半关闭心跳;RST/端点销毁 +读错误;FIN 先于 accept、FIN 先于 connect 返回;各握手阶段的心跳 +筛选。通知饱和测试在旧实现第五次 accept 超时;提前 FIN 测试在修复 +前等待 EOF 超时。排空测试使用多工作线程运行时,并做过额外重复验证。 + +主要命令(Linux 容器内,Rust 1.95): + +```sh +cargo +1.95 test --locked -p easytier-core \ + --features proxy-smoltcp-stack --lib gateway:: +cargo +1.95 nextest run --locked -p easytier --features full --lib \ + -E 'test(subnet_proxy_three_node_test) | test(subnet_proxy_half_close_test) | test(acl_rule_test_inbound) | test(acl_rule_test_subnet_proxy) | test(proxy_three_node_disconnect_test) | test(config_patch_test) | test(port_forward_with_inbound_default_drop_acl_test)' \ + --test-threads 1 --no-fail-fast +cargo fmt --all -- --check +cargo +1.95 clippy --locked -p easytier --features full \ + --lib --tests -- -D warnings +``` + +网关严格 Clippy 另以 `proxy-smoltcp-stack,ring-crypto` 通过;原因见上一轮 +记录中的既有未使用警告说明。Windows 所选 stable 工具链缺少格式检查 +组件,未将 Windows 格式或 Clippy 计为通过。 + +## 最终二进制真实流量 + +普通 TCP、KCP、QUIC 与内核/smoltcp 六种组合,底层 UDP;每组合 300 次 +短连接、16 并发、5 秒 socket 超时,检查 greeting 和逐字节回显。 + +- 六模式主矩阵:1,800/1,800 次短连接通过。 +- 请求 EOF 后返回响应:小/大请求共 12/12,通过;大请求 256 KiB, + 响应 1 MiB + 4 字节。 +- 服务端先关闭写端,客户端读到 EOF 后再发 256 KiB:6/6,通过, + 服务端实际校验接收数据。 +- connect 后立即关闭写端,不发送模式字节、不等待 greeting:每模式 + 32 次、16 并发,共 192/192 次纯空请求通过。 +- KCP 每栈额外三轮,与主矩阵合计 2,400/2,400 次短连接通过。追加 + 测试的大小半关闭及 192 次纯空请求通过,6 次反向半关闭中 1 次超时, + 因此不能将所有追加半关闭计为通过,具体证据见下文。 +- 最终版本在 `delay 10ms 3ms loss 1%` 下,六模式各两个连接持续 + 15 秒逐字节回显、12 次大小半关闭及 6 次反向半关闭通过。KCP 短连接 + 仍有内核模式 5/300、smoltcp 模式 8/300 次 greeting 超时,其他 + 四种模式零失败;没有为每次超时单独抓包证明原因相同。 +- 新旧两方向 × 六模式,共 12 组合正常持续回显通过。旧源端→新目标端 + 1,800 次短连接全通过;新源端→旧目标端 KCP 内核模式 12/300、 + smoltcp 模式 9/300 次超时,其他四种模式零失败。 + +旧版同样的无丢包 KCP 2,400 次短连接有 78 次 greeting 超时;旧版六种 +模式在请求 EOF 后均没有响应,且出现过 QUIC 请求数据被截断。 + +最终版本 KCP 主矩阵后,两种栈源端各 335 条 Closed 代理记录均在 +15 秒观察点清零;目标端记录为零。两端 FD 数保持内核模式 15、smoltcp +模式 14。RSS 较起始保留约 3.2–4.4 MiB 增长,单次有限负载不能证明 +或排除长期内存泄漏,也不能将代理条目清零等同于所有内部对象回收。 + +## 仍然存在的丢包握手问题与升级限制 + +在源端施加 `netem delay 10ms 3ms loss 1%` 后,中间修复版本 +(SHA-256 `2e78ecf…`)KCP 仍有 3/600 次 greeting 超时。额外有界 +TRACE 重跑捕获 3 个连接:源端收到 SYN|ACK,发出空 ACK|DATA 后即 +返回已建立;目标端未收到最终 ACK,直到 5 秒后的 FIN 才返回 RST。 +这与 accept 通知饱和不同,目标端尚未完成握手。 + +该控制报文可靠性问题需要单独设计握手恢复、重传及过期清理的一致性, +不能仅延长超时或重试应用请求。本次没有声称消除了所有 KCP 超时。 +具体连接号和双方日志行见 `remaining-loss-handshake.md`、 +`loss-handshake-evidence.log`。 + +混合版本正常回显通过并不意味着旧端获得修复。旧 KCP 目标端仍保留 +accept 缺陷;KCP 的完整半关闭能力需要两端升级。没有修改报文格式。 + +追加无 netem 的 KCP/kernel 反向半关闭中,连接 `3969305677` 在等 +服务端数据和 EOF 时超时;服务端没有收到客户端原定在 EOF 后发送的 +256 KiB 数据。发送端同一连接出现 780 次数据输出队列 Full,随后在 +08:36:47.247 UTC 发送缓冲排空,FIN 成功进入输出队列;接收端直到 +5 秒业务超时都未报告对端关闭,08:37:01.030 才报告关闭或重置信号。 +该 debug 日志行也可能由 RST 触发,不能据此认定迟到的是 FIN,更不能 +将数据队列 Full 直接解释为 FIN 被该队列丢弃。现有 FIN 发送路径没有 +确认重传;这一单次失败的具体丢包点尚未证实,保留为未解决的验证异常。 +同场景有界 TRACE 重跑 3 轮,900 次短连接、3 次反向半关闭均通过, +抓包确认成功轮次收到 FIN 并完成反向数据,未复现原失败;不能以重跑 +通过覆盖原失败。旧依赖与最终依赖的发送排空后入队 FIN 路径相同, +对照保存在 `fin-send-path-old-new.txt`。 + +异常轮次的 15 秒资源观察点仍有 1 条 Closed 代理记录,无 Connecting/ +Connected,FD 已回到起始值;因此仅主矩阵可以报告该观察点全部清零。 +定向 TRACE 三轮的 1,005 条 Closed 记录在 15 秒观察点全部清零。 + +## 验证中的异常与证据 + +首轮未完成的集成运行中,端口转发 ACL case 3 曾在连接本地监听端口时 +出现一次 ConnectionRefused;随后完整一轮 276 项通过,同项也通过。 +旧版单次及额外十次均通过,因此没有证据把这次失败断言为旧版必现问题, +也未修改无关的端口转发代码。新增半关闭测试的早期夹具曾只等待路由 +公告而未等待实际 TUN 路由可用,已补足就绪检查,并用有界 try_join +及时报告连接错误;这不计入生产代码的 red/green 对照。 + +Windows 测试早期表现为挂起,原因是测试 join 在 connect 失败后仍等待 +accept。限定握手等待后暴露出上述 PING/RST/SYN 竞态,旧依赖和修复前 +分支均有 TRACE 证据;最终原生测试通过。完整 diff 和后续心跳增量均 +经过独立子代理审查,没有 high-confidence 的新缺陷发现。 + +原始脚本、日志、逐连接 JSON、二进制和哈希保存在本机: + +```text +/data/project/proxy-close-validation-20260914/ +``` + +主要文件:`delivery-manifest.json`、`delivery-integration.log`、`delivery-clippy.log`、 +`final-core-gateway.log`、`macos-core.log`、`windows-core.log`、 +`mac-kcp-heartbeat-tests.log`、`final-tests.log`(Windows KCP)、 +`delivery-summary.json`、`delivery-kcp-repeat-summary.json`、 +`delivery-*-resources.json`、`traffic-summary.md`、`traffic-binaries.json`、 +`delivery-commands.json`、`artifacts-manifest.json`、`delivery-reverse-timeout.md`、 +`relay-red.log`、`windows-kcp-trace-fail.txt`、 +`windows-kcp-baseline-startup.txt`。Windows 部分原始日志为 UTF-16LE。 +Linux kcp-sys 测试、格式、Clippy 及依赖测试 red/green 对照为子代理 +执行结果,未单独落盘原始命令日志;最终 10 项测试耗时 11.03 秒。 + +未覆盖 Windows/macOS 原生 TUN 端到端路径、生产规模长时间运行、吞吐 +基准及容量极限负载。流量实验仅使用并清理自己的 namespace 和进程。 diff --git a/easytier-contrib/easytier-ohrs/Cargo.lock b/easytier-contrib/easytier-ohrs/Cargo.lock index b43d241e..a5560c72 100644 --- a/easytier-contrib/easytier-ohrs/Cargo.lock +++ b/easytier-contrib/easytier-ohrs/Cargo.lock @@ -2573,7 +2573,7 @@ dependencies = [ [[package]] name = "kcp-sys" version = "0.1.0" -source = "git+https://github.com/EasyTier/kcp-sys?rev=d7427c22d764deb1860a7d37acc446ed5033464c#d7427c22d764deb1860a7d37acc446ed5033464c" +source = "git+https://github.com/EasyTier/kcp-sys?rev=268533568d734ae89dc89603078da3ca522effe1#268533568d734ae89dc89603078da3ca522effe1" dependencies = [ "anyhow", "auto_impl", diff --git a/easytier-core/src/gateway/proxy/tcp_proxy_engine.rs b/easytier-core/src/gateway/proxy/tcp_proxy_engine.rs index 5ea13d40..b4cd785d 100644 --- a/easytier-core/src/gateway/proxy/tcp_proxy_engine.rs +++ b/easytier-core/src/gateway/proxy/tcp_proxy_engine.rs @@ -9,7 +9,7 @@ use std::{ use cidr::Ipv4Inet; use crossbeam::atomic::AtomicCell; -use dashmap::DashMap; +use dashmap::{DashMap, mapref::entry::Entry}; use smoltcp::wire::{IpAddress, IpProtocol, Ipv4Packet, TcpPacket}; use crate::packet::{PacketType, ZCPacket}; @@ -25,6 +25,16 @@ pub(crate) enum TcpProxyMode { QuicSrc, } +impl TcpProxyMode { + pub(super) const fn smoltcp_listener_port(self) -> u16 { + match self { + Self::Tcp => 8899, + Self::KcpSrc => 8900, + Self::QuicSrc => 8901, + } + } +} + #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum TcpNatEntryState { SynReceived, @@ -44,24 +54,30 @@ pub struct TcpNatEntrySnapshot { pub state: TcpNatEntryState, } +#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] +struct TcpNatFlowKey { + src: SocketAddr, + mapped_dst: SocketAddr, +} + #[derive(Debug)] pub(crate) struct TcpNatEntry { id: TcpNatEntryId, - src: SocketAddr, + flow: TcpNatFlowKey, + translated_src: SocketAddr, real_dst: SocketAddr, - mapped_dst: SocketAddr, start_time: Instant, start_time_unix_secs: u64, state: AtomicCell, } impl TcpNatEntry { - fn new(src: SocketAddr, real_dst: SocketAddr, mapped_dst: SocketAddr) -> Self { + fn new(flow: TcpNatFlowKey, translated_src: SocketAddr, real_dst: SocketAddr) -> Self { Self { id: uuid::Uuid::new_v4(), - src, + flow, + translated_src, real_dst, - mapped_dst, start_time: Instant::now(), start_time_unix_secs: SystemTime::now() .duration_since(UNIX_EPOCH) @@ -76,7 +92,7 @@ impl TcpNatEntry { } pub fn src(&self) -> SocketAddr { - self.src + self.flow.src } pub fn real_dst(&self) -> SocketAddr { @@ -84,7 +100,7 @@ impl TcpNatEntry { } pub fn mapped_dst(&self) -> SocketAddr { - self.mapped_dst + self.flow.mapped_dst } pub fn state(&self) -> TcpNatEntryState { @@ -97,9 +113,9 @@ impl TcpNatEntry { fn snapshot(&self) -> TcpNatEntrySnapshot { TcpNatEntrySnapshot { - src: self.src, + src: self.src(), dst: self.real_dst, - mapped_dst: self.mapped_dst, + mapped_dst: self.mapped_dst(), start_time: self.start_time_unix_secs, state: self.state(), } @@ -127,6 +143,7 @@ pub(crate) struct TcpProxyNicContext { #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub(crate) enum TcpProxyPacketAction { Handled { new_syn: bool }, + Drop, Pass, } @@ -134,9 +151,10 @@ pub(crate) enum TcpProxyPacketAction { pub(crate) struct TcpProxyEngine { cidr_table: Arc, local_port: AtomicU16, - syn_map: DashMap>, + next_translated_port: AtomicU16, + flow_map: DashMap>, + translated_src_map: DashMap>, conn_map: DashMap>, - addr_conn_map: DashMap>, } impl TcpProxyEngine { @@ -144,9 +162,10 @@ impl TcpProxyEngine { Self { cidr_table, local_port: AtomicU16::new(0), - syn_map: DashMap::new(), + next_translated_port: AtomicU16::new(1), + flow_map: DashMap::new(), + translated_src_map: DashMap::new(), conn_map: DashMap::new(), - addr_conn_map: DashMap::new(), } } @@ -158,6 +177,61 @@ impl TcpProxyEngine { self.local_port.load(Ordering::Relaxed) } + fn allocate_entry( + &self, + flow: TcpNatFlowKey, + real_dst: SocketAddr, + ) -> Option> { + let local_port = self.local_port(); + for _ in 0..u16::MAX { + let translated_port = self + .next_translated_port + .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |port| { + Some(if port == u16::MAX { 1 } else { port + 1 }) + }) + .expect("translated port counter update cannot fail"); + if translated_port == local_port { + continue; + } + let translated_src = SocketAddr::new(flow.src.ip(), translated_port); + let Entry::Vacant(slot) = self.translated_src_map.entry(translated_src) else { + continue; + }; + let entry = Arc::new(TcpNatEntry::new(flow, translated_src, real_dst)); + slot.insert(entry.clone()); + return Some(entry); + } + None + } + + fn entry_for_syn( + &self, + flow: TcpNatFlowKey, + real_dst: SocketAddr, + ) -> Option<(Arc, bool)> { + match self.flow_map.entry(flow) { + Entry::Occupied(mut slot) => { + let entry = slot.get(); + if !matches!( + entry.state(), + TcpNatEntryState::ClosingSrc + | TcpNatEntryState::ClosingDst + | TcpNatEntryState::Closed + ) { + return Some((entry.clone(), false)); + } + let entry = self.allocate_entry(flow, real_dst)?; + slot.insert(entry.clone()); + Some((entry, true)) + } + Entry::Vacant(slot) => { + let entry = self.allocate_entry(flow, real_dst)?; + slot.insert(entry.clone()); + Some((entry, true)) + } + } + } + pub fn check_packet_from_peer_fast( &self, mode: TcpProxyMode, @@ -230,28 +304,36 @@ impl TcpProxyEngine { let source_ip = ip_packet.src_addr(); let source_port = tcp_packet.src_port(); - let src = SocketAddr::V4(SocketAddrV4::new(source_ip, source_port)); + let dest_ip = ip_packet.dst_addr(); + let dest_port = tcp_packet.dst_port(); + let flow = TcpNatFlowKey { + src: SocketAddr::V4(SocketAddrV4::new(source_ip, source_port)), + mapped_dst: SocketAddr::V4(SocketAddrV4::new(dest_ip, dest_port)), + }; - let mut new_syn = false; - if tcp_packet.syn() && !tcp_packet.ack() { - let dest_ip = ip_packet.dst_addr(); - let dest_port = tcp_packet.dst_port(); - let mapped_dst = SocketAddr::V4(SocketAddrV4::new(dest_ip, dest_port)); + let is_syn = tcp_packet.syn() && !tcp_packet.ack(); + let (entry, new_syn) = if is_syn { let real_dst = SocketAddr::V4(SocketAddrV4::new(real_dst_ip, dest_port)); + let Some(entry) = self.entry_for_syn(flow, real_dst) else { + tracing::error!(?flow, "tcp proxy translated source ports exhausted"); + return TcpProxyPacketAction::Drop; + }; + entry + } else { + let Some(entry) = self.flow_map.get(&flow) else { + return TcpProxyPacketAction::Pass; + }; + (entry.clone(), false) + }; - let old_val = self - .syn_map - .insert(src, Arc::new(TcpNatEntry::new(src, real_dst, mapped_dst))); + if new_syn { tracing::info!( - ?src, - ?real_dst, - ?mapped_dst, - old_entry = ?old_val, + src = ?entry.src(), + translated_src = ?entry.translated_src, + real_dst = ?entry.real_dst(), + mapped_dst = ?entry.mapped_dst(), "tcp syn received" ); - new_syn = true; - } else if !self.addr_conn_map.contains_key(&src) && !self.syn_map.contains_key(&src) { - return TcpProxyPacketAction::Pass; } let mut ip_packet = Ipv4Packet::new_checked(payload_bytes).expect("checked ipv4 packet"); @@ -263,6 +345,7 @@ impl TcpProxyEngine { { let mut tcp_packet = TcpPacket::new_checked(ip_packet.payload_mut()).expect("checked tcp packet"); + tcp_packet.set_src_port(entry.translated_src.port()); tcp_packet.set_dst_port(ctx.local_port); tcp_packet.fill_checksum(&IpAddress::Ipv4(source), &IpAddress::Ipv4(local_ip)); } @@ -312,17 +395,13 @@ impl TcpProxyEngine { } tracing::trace!(?dst_addr, "tcp packet try find entry"); - let entry = if let Some(entry) = self.addr_conn_map.get(&dst_addr) { - entry.clone() - } else { - let Some(syn_entry) = self.syn_map.get(&dst_addr) else { - return false; - }; - syn_entry.clone() + let Some(entry) = self.translated_src_map.get(&dst_addr) else { + return false; }; - assert_eq!(entry.src, dst_addr); + let entry = entry.clone(); + assert_eq!(entry.translated_src, dst_addr); - let IpAddr::V4(mapped_dst_ip) = entry.mapped_dst.ip() else { + let IpAddr::V4(mapped_dst_ip) = entry.mapped_dst().ip() else { panic!("v4 nat entry src ip is not v4"); }; @@ -346,6 +425,7 @@ impl TcpProxyEngine { let mut tcp_packet = TcpPacket::new_checked(ip_packet.payload_mut()).expect("checked tcp packet"); tcp_packet.set_src_port(entry.real_dst.port()); + tcp_packet.set_dst_port(entry.src().port()); tcp_packet.fill_checksum(&IpAddress::Ipv4(mapped_dst_ip), &IpAddress::Ipv4(dst)); } ip_packet.fill_checksum(); @@ -366,53 +446,90 @@ impl TcpProxyEngine { } } - let (_, entry) = self.syn_map.remove(&socket_addr)?; - if entry.state() != TcpNatEntryState::SynReceived { + let entry = self.translated_src_map.get(&socket_addr)?.clone(); + if entry + .state + .compare_exchange( + TcpNatEntryState::SynReceived, + TcpNatEntryState::ConnectingDst, + ) + .is_err() + { + if entry.state() == TcpNatEntryState::Closed { + self.remove_indices(&entry); + } return None; } - entry.set_state(TcpNatEntryState::ConnectingDst); - self.addr_conn_map.insert(entry.src, entry.clone()); let old_nat_val = self.conn_map.insert(entry.id, entry.clone()); assert!(old_nat_val.is_none()); Some(entry) } + fn remove_indices(&self, entry: &TcpNatEntry) { + self.flow_map + .remove_if(&entry.flow, |_, current| current.id == entry.id); + self.translated_src_map + .remove_if(&entry.translated_src, |_, current| current.id == entry.id); + } + pub fn remove_entry(&self, entry_id: TcpNatEntryId) { let Some((_, entry)) = self.conn_map.remove(&entry_id) else { return; }; - self.addr_conn_map - .remove_if(&entry.src, |_, current| current.id == entry.id); + self.remove_indices(&entry); if self.conn_map.capacity() - self.conn_map.len() > 16 { self.conn_map.shrink_to_fit(); } - if self.addr_conn_map.capacity() - self.addr_conn_map.len() > 16 { - self.addr_conn_map.shrink_to_fit(); + if self.flow_map.capacity() - self.flow_map.len() > 16 { + self.flow_map.shrink_to_fit(); + } + if self.translated_src_map.capacity() - self.translated_src_map.len() > 16 { + self.translated_src_map.shrink_to_fit(); } } - pub fn cleanup_expired_syn(&self, timeout: Duration) { - self.syn_map.retain(|_, entry| { - if entry.start_time.elapsed() > timeout { - tracing::warn!(?entry, "syn nat entry expired"); - entry.set_state(TcpNatEntryState::Closed); - false - } else { - true - } - }); - self.syn_map.shrink_to_fit(); + pub fn clear(&self) { + for entry in self.flow_map.iter() { + entry.set_state(TcpNatEntryState::Closed); + } + for entry in self.conn_map.iter() { + entry.set_state(TcpNatEntryState::Closed); + } + self.flow_map.clear(); + self.translated_src_map.clear(); + self.conn_map.clear(); } - pub fn is_tcp_proxy_connection(&self, src: SocketAddr) -> bool { - self.syn_map.contains_key(&src) || self.addr_conn_map.contains_key(&src) + pub fn cleanup_expired_syn(&self, timeout: Duration) { + self.flow_map.retain(|_, entry| { + let expired = entry.start_time.elapsed() > timeout + && entry + .state + .compare_exchange(TcpNatEntryState::SynReceived, TcpNatEntryState::Closed) + .is_ok(); + if expired { + tracing::warn!(?entry, "syn nat entry expired"); + self.translated_src_map + .remove_if(&entry.translated_src, |_, current| current.id == entry.id); + } + !expired + }); + self.flow_map.shrink_to_fit(); + self.translated_src_map.shrink_to_fit(); + } + + pub fn is_tcp_proxy_flow(&self, src: SocketAddr, mapped_dst: SocketAddr) -> bool { + self.flow_map + .contains_key(&TcpNatFlowKey { src, mapped_dst }) } pub fn list_entries(&self) -> Vec { let mut entries = Vec::new(); - for entry in self.syn_map.iter() { - entries.push(entry.value().snapshot()); + for entry in self.flow_map.iter() { + if entry.state() == TcpNatEntryState::SynReceived { + entries.push(entry.value().snapshot()); + } } for entry in self.conn_map.iter() { entries.push(entry.value().snapshot()); @@ -503,6 +620,32 @@ mod tests { } } + fn packet_src(packet: &ZCPacket) -> SocketAddrV4 { + let ipv4 = Ipv4Packet::new_checked(packet.payload()).unwrap(); + let tcp = TcpPacket::new_checked(ipv4.payload()).unwrap(); + SocketAddrV4::new(ipv4.src_addr(), tcp.src_port()) + } + + fn nic_ctx() -> TcpProxyNicContext { + TcpProxyNicContext { + local_inet: Some("10.144.144.204/24".parse().unwrap()), + local_port: 8899, + my_peer_id: 2, + smoltcp_enabled: false, + } + } + + #[test] + fn smoltcp_listener_ports_are_unique_per_proxy_mode() { + let tcp = TcpProxyMode::Tcp.smoltcp_listener_port(); + let kcp = TcpProxyMode::KcpSrc.smoltcp_listener_port(); + let quic = TcpProxyMode::QuicSrc.smoltcp_listener_port(); + + assert_ne!(tcp, kcp); + assert_ne!(tcp, quic); + assert_ne!(kcp, quic); + } + #[test] fn peer_syn_creates_entry_and_rewrites_to_local_stack() { let engine = tcp_engine(); @@ -522,7 +665,7 @@ mod tests { "10.144.144.204".parse::().unwrap() ); let tcp = TcpPacket::new_checked(ipv4.payload()).unwrap(); - assert_eq!(tcp.src_port(), src.port()); + assert_ne!(tcp.src_port(), src.port()); assert_eq!(tcp.dst_port(), 8899); let entries = engine.list_entries(); @@ -545,25 +688,18 @@ mod tests { engine.try_handle_peer_packet(TcpProxyMode::Tcp, &mut request, peer_ctx()), TcpProxyPacketAction::Handled { new_syn: true } )); + let translated_src = packet_src(&request); let entry = engine .accept_connection( - SocketAddr::V4(src), + SocketAddr::V4(translated_src), Some("10.144.144.204/24".parse().unwrap()), ) .unwrap(); assert_eq!(entry.state(), TcpNatEntryState::ConnectingDst); let local = SocketAddrV4::new("10.144.144.204".parse().unwrap(), 8899); - let mut response = build_tcp_packet(local, src, false, true); - assert!(engine.try_process_packet_from_nic( - &mut response, - TcpProxyNicContext { - local_inet: Some("10.144.144.204/24".parse().unwrap()), - local_port: 8899, - my_peer_id: 2, - smoltcp_enabled: false, - }, - )); + let mut response = build_tcp_packet(local, translated_src, false, true); + assert!(engine.try_process_packet_from_nic(&mut response, nic_ctx())); let hdr: &PeerManagerHeader = response.peer_manager_header().unwrap(); assert!(hdr.is_no_proxy()); @@ -573,6 +709,111 @@ mod tests { let tcp = TcpPacket::new_checked(ipv4.payload()).unwrap(); assert_eq!(tcp.src_port(), mapped_dst.port()); assert_eq!(tcp.dst_port(), src.port()); + + engine.remove_entry(entry.id()); + assert!(!engine.is_tcp_proxy_flow(SocketAddr::V4(src), SocketAddr::V4(mapped_dst),)); + assert!( + engine + .translated_src_map + .get(&SocketAddr::V4(translated_src)) + .is_none() + ); + } + + #[test] + fn same_source_port_to_mapped_and_real_destinations_stay_distinct() { + let engine = tcp_engine(); + let src = SocketAddrV4::new("10.144.144.206".parse().unwrap(), 50000); + let mapped_dst = SocketAddrV4::new("10.10.10.42".parse().unwrap(), 80); + let real_dst = SocketAddrV4::new("127.0.0.42".parse().unwrap(), 80); + + let mut mapped_request = build_tcp_packet(src, mapped_dst, true, false); + assert_eq!( + engine.try_handle_peer_packet(TcpProxyMode::Tcp, &mut mapped_request, peer_ctx()), + TcpProxyPacketAction::Handled { new_syn: true } + ); + let mapped_translated_src = packet_src(&mapped_request); + engine + .accept_connection( + SocketAddr::V4(mapped_translated_src), + Some("10.144.144.204/24".parse().unwrap()), + ) + .unwrap(); + + let mut real_request = build_tcp_packet(src, real_dst, true, false); + real_request + .mut_peer_manager_header() + .unwrap() + .set_exit_node(true); + let mut real_ctx = peer_ctx(); + real_ctx.enable_exit_node = true; + assert_eq!( + engine.try_handle_peer_packet(TcpProxyMode::Tcp, &mut real_request, real_ctx), + TcpProxyPacketAction::Handled { new_syn: true } + ); + let real_translated_src = packet_src(&real_request); + engine + .accept_connection( + SocketAddr::V4(real_translated_src), + Some("10.144.144.204/24".parse().unwrap()), + ) + .unwrap(); + + assert_ne!(mapped_translated_src, real_translated_src); + + let local = SocketAddrV4::new("10.144.144.204".parse().unwrap(), 8899); + for (translated_src, expected_source) in [ + (mapped_translated_src, mapped_dst), + (real_translated_src, real_dst), + ] { + let mut response = build_tcp_packet(local, translated_src, false, true); + assert!(engine.try_process_packet_from_nic(&mut response, nic_ctx())); + + let ipv4 = Ipv4Packet::new_checked(response.payload()).unwrap(); + assert_eq!(ipv4.src_addr(), *expected_source.ip()); + assert_eq!(ipv4.dst_addr(), *src.ip()); + let tcp = TcpPacket::new_checked(ipv4.payload()).unwrap(); + assert_eq!(tcp.src_port(), expected_source.port()); + assert_eq!(tcp.dst_port(), src.port()); + } + } + + #[test] + fn clear_discards_accepted_entries_before_restart() { + let engine = tcp_engine(); + let src = SocketAddrV4::new("10.144.144.206".parse().unwrap(), 50000); + let mapped_dst = SocketAddrV4::new("10.10.10.42".parse().unwrap(), 80); + let mut request = build_tcp_packet(src, mapped_dst, true, false); + assert_eq!( + engine.try_handle_peer_packet(TcpProxyMode::Tcp, &mut request, peer_ctx()), + TcpProxyPacketAction::Handled { new_syn: true } + ); + let old_entry = engine + .accept_connection( + SocketAddr::V4(packet_src(&request)), + Some("10.144.144.204/24".parse().unwrap()), + ) + .unwrap(); + + engine.clear(); + + assert_eq!(old_entry.state(), TcpNatEntryState::Closed); + assert!(engine.flow_map.is_empty()); + assert!(engine.translated_src_map.is_empty()); + assert!(engine.conn_map.is_empty()); + + let mut retry = build_tcp_packet(src, mapped_dst, true, false); + assert_eq!( + engine.try_handle_peer_packet(TcpProxyMode::Tcp, &mut retry, peer_ctx()), + TcpProxyPacketAction::Handled { new_syn: true } + ); + let new_entry = engine + .accept_connection( + SocketAddr::V4(packet_src(&retry)), + Some("10.144.144.204/24".parse().unwrap()), + ) + .unwrap(); + assert_ne!(old_entry.id(), new_entry.id()); } #[test] @@ -585,19 +826,264 @@ mod tests { engine.try_handle_peer_packet(TcpProxyMode::Tcp, &mut request, peer_ctx()), TcpProxyPacketAction::Handled { new_syn: true } )); - let entry = engine.syn_map.get(&SocketAddr::V4(src)).unwrap().clone(); + let translated_src = packet_src(&request); + let entry = engine + .translated_src_map + .get(&SocketAddr::V4(translated_src)) + .unwrap() + .clone(); entry.set_state(TcpNatEntryState::Closed); assert!( engine .accept_connection( - SocketAddr::V4(src), + SocketAddr::V4(translated_src), Some("10.144.144.204/24".parse().unwrap()), ) .is_none() ); - assert!(engine.syn_map.get(&SocketAddr::V4(src)).is_none()); - assert!(engine.addr_conn_map.get(&SocketAddr::V4(src)).is_none()); + assert!(engine.flow_map.is_empty()); + assert!(engine.translated_src_map.is_empty()); assert!(engine.conn_map.is_empty()); } + + #[test] + fn retransmitted_syn_keeps_the_pending_and_accepted_mapping() { + let engine = tcp_engine(); + let src = "10.144.144.206:50000".parse().unwrap(); + let dst = "10.10.10.42:80".parse().unwrap(); + let mut first = build_tcp_packet(src, dst, true, false); + assert_eq!( + engine.try_handle_peer_packet(TcpProxyMode::Tcp, &mut first, peer_ctx()), + TcpProxyPacketAction::Handled { new_syn: true } + ); + let translated_src = SocketAddr::V4(packet_src(&first)); + + for accepted in [false, true] { + if accepted { + engine.accept_connection(translated_src, None).unwrap(); + engine.cleanup_expired_syn(Duration::ZERO); + } + let mut retry = build_tcp_packet(src, dst, true, false); + assert_eq!( + engine.try_handle_peer_packet(TcpProxyMode::Tcp, &mut retry, peer_ctx()), + TcpProxyPacketAction::Handled { new_syn: false } + ); + assert_eq!(SocketAddr::V4(packet_src(&retry)), translated_src); + assert_eq!(engine.flow_map.len(), 1); + assert_eq!(engine.translated_src_map.len(), 1); + assert_eq!(engine.list_entries().len(), 1); + } + assert!(engine.accept_connection(translated_src, None).is_none()); + } + + #[test] + fn removing_closing_connection_preserves_replacement_flow() { + for state in [ + TcpNatEntryState::ClosingSrc, + TcpNatEntryState::ClosingDst, + TcpNatEntryState::Closed, + ] { + let engine = tcp_engine(); + let flow = TcpNatFlowKey { + src: "10.144.144.206:50000".parse().unwrap(), + mapped_dst: "10.10.10.42:80".parse().unwrap(), + }; + let real_dst = "127.0.0.42:80".parse().unwrap(); + let (old, _) = engine.entry_for_syn(flow, real_dst).unwrap(); + engine.accept_connection(old.translated_src, None).unwrap(); + old.set_state(state); + let (replacement, new_syn) = engine.entry_for_syn(flow, real_dst).unwrap(); + assert!(new_syn); + assert_ne!(old.translated_src, replacement.translated_src); + assert_eq!(engine.translated_src_map.len(), 2); + + engine.remove_entry(old.id()); + assert!(engine.is_tcp_proxy_flow(flow.src, flow.mapped_dst)); + assert!(!engine.translated_src_map.contains_key(&old.translated_src)); + assert_eq!(engine.flow_map.get(&flow).unwrap().id(), replacement.id()); + assert_eq!( + engine + .accept_connection(replacement.translated_src, None) + .unwrap() + .id(), + replacement.id() + ); + + let local = "10.144.144.204:8899".parse().unwrap(); + let SocketAddr::V4(translated_src) = replacement.translated_src else { + unreachable!(); + }; + let mut response = build_tcp_packet(local, translated_src, false, true); + assert!(engine.try_process_packet_from_nic(&mut response, nic_ctx())); + let ip = Ipv4Packet::new_checked(response.payload()).unwrap(); + let tcp = TcpPacket::new_checked(ip.payload()).unwrap(); + assert_eq!(SocketAddr::V4(packet_src(&response)), flow.mapped_dst); + assert_eq!(tcp.dst_port(), flow.src.port()); + assert!(ip.verify_checksum()); + assert!(tcp.verify_checksum( + &IpAddress::Ipv4(ip.src_addr()), + &IpAddress::Ipv4(ip.dst_addr()) + )); + } + } + + #[test] + fn expired_syn_releases_both_indices_without_removing_accepted_flow() { + let engine = tcp_engine(); + let pending_flow = TcpNatFlowKey { + src: "10.144.144.206:50000".parse().unwrap(), + mapped_dst: "10.10.10.42:80".parse().unwrap(), + }; + let accepted_flow = TcpNatFlowKey { + mapped_dst: "10.10.10.43:80".parse().unwrap(), + ..pending_flow + }; + let real_dst = "127.0.0.42:80".parse().unwrap(); + let (pending, _) = engine.entry_for_syn(pending_flow, real_dst).unwrap(); + let (accepted, _) = engine.entry_for_syn(accepted_flow, real_dst).unwrap(); + engine + .accept_connection(accepted.translated_src, None) + .unwrap(); + + engine.cleanup_expired_syn(Duration::ZERO); + + assert_eq!(pending.state(), TcpNatEntryState::Closed); + assert!( + engine + .accept_connection(pending.translated_src, None) + .is_none() + ); + assert!(!engine.is_tcp_proxy_flow(pending_flow.src, pending_flow.mapped_dst)); + assert!(engine.is_tcp_proxy_flow(accepted_flow.src, accepted_flow.mapped_dst)); + assert_eq!(engine.translated_src_map.len(), 1); + assert_eq!(engine.conn_map.len(), 1); + assert_eq!(accepted.state(), TcpNatEntryState::ConnectingDst); + } + + #[test] + fn translated_port_wrap_skips_live_ports_and_listener() { + let engine = tcp_engine(); + engine.set_local_port(8899); + let flow = TcpNatFlowKey { + src: "10.144.144.206:50000".parse().unwrap(), + mapped_dst: "10.10.10.42:80".parse().unwrap(), + }; + let real_dst = "127.0.0.42:80".parse().unwrap(); + engine + .next_translated_port + .store(u16::MAX, Ordering::Relaxed); + let last = engine.allocate_entry(flow, real_dst).unwrap(); + let first = engine.allocate_entry(flow, real_dst).unwrap(); + assert_eq!(last.translated_src.port(), u16::MAX); + assert_eq!(first.translated_src.port(), 1); + + engine + .next_translated_port + .store(u16::MAX, Ordering::Relaxed); + let next = engine.allocate_entry(flow, real_dst).unwrap(); + assert_eq!(next.translated_src.port(), 2); + assert_eq!( + engine + .translated_src_map + .get(&last.translated_src) + .unwrap() + .id(), + last.id() + ); + assert_eq!( + engine + .translated_src_map + .get(&first.translated_src) + .unwrap() + .id(), + first.id() + ); + + engine.next_translated_port.store(8899, Ordering::Relaxed); + let after_listener = engine.allocate_entry(flow, real_dst).unwrap(); + assert_eq!(after_listener.translated_src.port(), 8900); + } + + #[test] + fn exhausted_translated_ports_are_isolated_and_reusable_after_cleanup() { + let engine = tcp_engine(); + engine.set_local_port(8899); + let src_ip = "10.144.144.206".parse().unwrap(); + let dst = "10.10.10.42:80".parse().unwrap(); + let real_dst = "127.0.0.42:80".parse().unwrap(); + for port in 1..u16::MAX { + let flow = TcpNatFlowKey { + src: SocketAddr::V4(SocketAddrV4::new(src_ip, port)), + mapped_dst: SocketAddr::V4(dst), + }; + assert!(engine.entry_for_syn(flow, real_dst).is_some()); + } + assert_eq!(engine.translated_src_map.len(), usize::from(u16::MAX) - 1); + let src = SocketAddrV4::new(src_ip, u16::MAX); + let mut exhausted = build_tcp_packet(src, dst, true, false); + assert_eq!( + engine.try_handle_peer_packet(TcpProxyMode::Tcp, &mut exhausted, peer_ctx()), + TcpProxyPacketAction::Drop + ); + assert!(!engine.is_tcp_proxy_flow(SocketAddr::V4(src), SocketAddr::V4(dst))); + + let other_flow = TcpNatFlowKey { + src: "10.144.144.207:50000".parse().unwrap(), + mapped_dst: SocketAddr::V4(dst), + }; + assert!(engine.entry_for_syn(other_flow, real_dst).is_some()); + + let released_addr = SocketAddr::V4(SocketAddrV4::new(src_ip, 1)); + let released = engine.accept_connection(released_addr, None).unwrap(); + released.set_state(TcpNatEntryState::Closed); + engine.remove_entry(released.id()); + let mut retry = build_tcp_packet(src, dst, true, false); + assert_eq!( + engine.try_handle_peer_packet(TcpProxyMode::Tcp, &mut retry, peer_ctx()), + TcpProxyPacketAction::Handled { new_syn: true } + ); + assert_eq!(SocketAddr::V4(packet_src(&retry)), released_addr); + assert!(engine.accept_connection(released_addr, None).is_some()); + + engine.clear(); + assert!(engine.flow_map.is_empty()); + assert!(engine.translated_src_map.is_empty()); + assert!(engine.conn_map.is_empty()); + assert!(engine.entry_for_syn(other_flow, real_dst).is_some()); + } + + #[test] + fn concurrent_syns_share_one_mapping_and_only_one_accept() { + let engine = tcp_engine(); + let barrier = std::sync::Barrier::new(8); + let flow = TcpNatFlowKey { + src: "10.144.144.206:50000".parse().unwrap(), + mapped_dst: "10.10.10.42:80".parse().unwrap(), + }; + let results = std::thread::scope(|scope| { + let tasks: Vec<_> = (0..8) + .map(|_| { + scope.spawn(|| { + barrier.wait(); + let (entry, new_syn) = engine.entry_for_syn(flow, flow.mapped_dst).unwrap(); + let accepted = engine + .accept_connection(entry.translated_src, None) + .is_some(); + (entry.id(), new_syn, accepted) + }) + }) + .collect(); + tasks + .into_iter() + .map(|task| task.join().unwrap()) + .collect::>() + }); + assert!(results.iter().all(|result| result.0 == results[0].0)); + assert_eq!(results.iter().filter(|result| result.1).count(), 1); + assert_eq!(results.iter().filter(|result| result.2).count(), 1); + assert_eq!(engine.flow_map.len(), 1); + assert_eq!(engine.translated_src_map.len(), 1); + assert_eq!(engine.conn_map.len(), 1); + } } diff --git a/easytier-core/src/gateway/proxy/tcp_proxy_service.rs b/easytier-core/src/gateway/proxy/tcp_proxy_service.rs index 223211f8..805bcb11 100644 --- a/easytier-core/src/gateway/proxy/tcp_proxy_service.rs +++ b/easytier-core/src/gateway/proxy/tcp_proxy_service.rs @@ -3,7 +3,8 @@ use std::sync::{Arc, Weak, atomic::Ordering}; use std::time::Duration; use atomic_shim::AtomicU64; -use tokio::io::{AsyncWriteExt, copy}; +use parking_lot::Mutex; +use tokio::io::AsyncWriteExt; use tokio::task::JoinSet; use crate::{ @@ -34,10 +35,10 @@ use crate::gateway::smoltcp::{SmolTcpStack, output_dst_ip}; fn spawn_tcp_proxy_task( lifecycle: &AtomicU64, expected_generation: u64, - tasks: &std::sync::Mutex>, + tasks: &Mutex>, task: impl Future + Send + 'static, ) -> bool { - let mut tasks = tasks.lock().unwrap(); + let mut tasks = tasks.lock(); if lifecycle.load(Ordering::Acquire) != expected_generation { return false; } @@ -57,12 +58,12 @@ pub struct TcpProxyService< connector: Arc, engine: Arc, mode: TcpProxyMode, - peer_pipeline_guard: std::sync::Mutex>, - nic_pipeline_guard: std::sync::Mutex>, - kernel_listener: std::sync::Mutex>>, + peer_pipeline_guard: Mutex>, + nic_pipeline_guard: Mutex>, + kernel_listener: Mutex>>, #[cfg(feature = "proxy-smoltcp-stack")] - smoltcp_stack: std::sync::Mutex>>, - tasks: std::sync::Mutex>, + smoltcp_stack: Mutex>>, + tasks: Mutex>, lifecycle: AtomicU64, } @@ -86,12 +87,12 @@ impl) { - if self.peer_pipeline_guard.lock().unwrap().is_some() { + if self.peer_pipeline_guard.lock().is_some() { return; } let peer_guard = self @@ -213,11 +218,11 @@ impl) { - if self.nic_pipeline_guard.lock().unwrap().is_some() { + if self.nic_pipeline_guard.lock().is_some() { return; } let nic_guard = self @@ -226,7 +231,7 @@ impl, generation: u64) { @@ -244,7 +249,7 @@ impl {} + TcpProxyPacketAction::Drop => return None, + TcpProxyPacketAction::Pass => return Some(packet), + } if snapshot.smoltcp_enabled { #[cfg(feature = "proxy-smoltcp-stack")] @@ -486,7 +490,7 @@ impl Result<(), ProxyRuntimeError> { - let (mut src_reader, mut src_writer) = tokio::io::split(src); - let (mut dst_reader, mut dst_writer) = tokio::io::split(dst); - let src_to_dst = copy(&mut src_reader, &mut dst_writer); - let dst_to_src = copy(&mut dst_reader, &mut src_writer); - tokio::pin!(src_to_dst); - tokio::pin!(dst_to_src); - tokio::select! { - result = &mut src_to_dst => { - result?; - } - result = &mut dst_to_src => { - result?; - } - } + // Forward EOF to the opposite writer while continuing to relay its reply. + // Returning on the first EOF would discard responses to half-closed requests. + tokio::io::copy_bidirectional(src, dst).await?; Ok(()) } @@ -575,6 +568,57 @@ impl); @@ -595,7 +639,7 @@ mod tests { #[tokio::test] async fn stop_fence_linearizes_task_registration() { let lifecycle = AtomicU64::new(1); - let tasks = std::sync::Mutex::new(JoinSet::new()); + let tasks = Mutex::new(JoinSet::new()); let accepted_dropped = Arc::new(AtomicBool::new(false)); assert!(spawn_tcp_proxy_task( @@ -605,7 +649,7 @@ mod tests { pending_task(accepted_dropped.clone()), )); lifecycle.store(2, Ordering::Release); - let mut stopping = std::mem::take(&mut *tasks.lock().unwrap()); + let mut stopping = std::mem::take(&mut *tasks.lock()); stopping.shutdown().await; assert!(accepted_dropped.load(Ordering::Acquire)); @@ -627,7 +671,7 @@ mod tests { &tasks, pending_task(current_dropped.clone()), )); - let mut current = std::mem::take(&mut *tasks.lock().unwrap()); + let mut current = std::mem::take(&mut *tasks.lock()); current.shutdown().await; assert!(current_dropped.load(Ordering::Acquire)); } diff --git a/easytier-core/src/gateway/proxy/wrapped_tcp_proxy.rs b/easytier-core/src/gateway/proxy/wrapped_tcp_proxy.rs index aa613c25..7c3090e2 100644 --- a/easytier-core/src/gateway/proxy/wrapped_tcp_proxy.rs +++ b/easytier-core/src/gateway/proxy/wrapped_tcp_proxy.rs @@ -128,11 +128,11 @@ pub struct WrappedTcpProxyNicContext { pub async fn try_process_wrapped_tcp_packet_from_nic( zc_packet: &mut ZCPacket, ctx: WrappedTcpProxyNicContext, - is_tcp_proxy_connection: ConnectionLookup, + is_tcp_proxy_flow: ConnectionLookup, check_dst_allowed: AllowCheck, ) -> bool where - ConnectionLookup: Fn(SocketAddr) -> bool, + ConnectionLookup: Fn(SocketAddr, SocketAddr) -> bool, AllowCheck: FnOnce(Ipv4Addr) -> AllowCheckFut, AllowCheckFut: Future, { @@ -156,6 +156,7 @@ where let src_ip = ip_packet.src_addr(); let dst_ip = ip_packet.dst_addr(); let src_port = tcp_packet.src_port(); + let dst_port = tcp_packet.dst_port(); let is_syn = tcp_packet.syn() && !tcp_packet.ack(); if is_syn { @@ -171,7 +172,10 @@ where ); return false; } - } else if !is_tcp_proxy_connection(SocketAddr::V4(SocketAddrV4::new(src_ip, src_port))) { + } else if !is_tcp_proxy_flow( + SocketAddr::V4(SocketAddrV4::new(src_ip, src_port)), + SocketAddr::V4(SocketAddrV4::new(dst_ip, dst_port)), + ) { return false; } @@ -366,7 +370,7 @@ mod tests { try_process_wrapped_tcp_packet_from_nic( &mut packet, context(WrappedTcpProxyTransport::Kcp), - |_| false, + |_, _| false, |_| async { true }, ) .await @@ -387,7 +391,7 @@ mod tests { !try_process_wrapped_tcp_packet_from_nic( &mut packet, context(WrappedTcpProxyTransport::Kcp), - |_| false, + |_, _| false, |_| async { false }, ) .await @@ -407,7 +411,9 @@ mod tests { try_process_wrapped_tcp_packet_from_nic( &mut packet, context(WrappedTcpProxyTransport::Quic), - |addr| addr == SocketAddr::V4(src), + |src_addr, dst_addr| { + src_addr == SocketAddr::V4(src) && dst_addr == SocketAddr::V4(dst) + }, |_| async { false }, ) .await @@ -428,7 +434,7 @@ mod tests { !try_process_wrapped_tcp_packet_from_nic( &mut packet, context(WrappedTcpProxyTransport::Quic), - |_| false, + |_, _| false, |_| async { true }, ) .await @@ -445,7 +451,7 @@ mod tests { !try_process_wrapped_tcp_packet_from_nic( &mut packet, context(WrappedTcpProxyTransport::Kcp), - |_| false, + |_, _| false, |_| async { true }, ) .await diff --git a/easytier-core/src/gateway/proxy/wrapped_transport/packet_plane.rs b/easytier-core/src/gateway/proxy/wrapped_transport/packet_plane.rs index c6664c28..d39eb676 100644 --- a/easytier-core/src/gateway/proxy/wrapped_transport/packet_plane.rs +++ b/easytier-core/src/gateway/proxy/wrapped_transport/packet_plane.rs @@ -260,7 +260,7 @@ where local_ipv4: snapshot.virtual_ipv4, smoltcp_enabled: snapshot.smoltcp_enabled, }, - move |src| connection_engine.is_tcp_proxy_connection(src), + move |src, mapped_dst| connection_engine.is_tcp_proxy_flow(src, mapped_dst), move |dst_ip| async move { match transport { WrappedTransportKind::Kcp => { diff --git a/easytier-core/src/gateway/smoltcp/stack.rs b/easytier-core/src/gateway/smoltcp/stack.rs index d20c2b33..84aa1622 100644 --- a/easytier-core/src/gateway/smoltcp/stack.rs +++ b/easytier-core/src/gateway/smoltcp/stack.rs @@ -11,6 +11,7 @@ use crate::gateway::proxy::traits::TcpProxyStream; pub struct SmolTcpStack { ingress_tx: mpsc::Sender, + local_port: u16, output_rx: Mutex>>>, listener: Mutex, _net: Net, @@ -18,7 +19,7 @@ pub struct SmolTcpStack { } impl SmolTcpStack { - pub async fn new(local_ip: Ipv4Addr) -> anyhow::Result> { + pub async fn new(local_ip: Ipv4Addr, local_port: u16) -> anyhow::Result> { let tasks = Arc::new(std::sync::Mutex::new(JoinSet::new())); let mut cap = smoltcp::phy::DeviceCapabilities::default(); cap.max_transmission_unit = 1280; @@ -63,12 +64,13 @@ impl SmolTcpStack { ); net.set_any_ip(true); let listener = net - .tcp_bind("0.0.0.0:8899".parse().unwrap()) + .tcp_bind(SocketAddr::new(Ipv4Addr::UNSPECIFIED.into(), local_port)) .await .map_err(|error| anyhow::anyhow!("bind smoltcp listener failed: {error}"))?; Ok(Arc::new(Self { ingress_tx, + local_port, output_rx: Mutex::new(Some(stack_stream)), listener: Mutex::new(listener), _net: net, @@ -77,7 +79,7 @@ impl SmolTcpStack { } pub fn local_port(&self) -> u16 { - 8899 + self.local_port } pub async fn send_ingress(&self, packet: ZCPacket) -> anyhow::Result<()> { @@ -137,12 +139,12 @@ mod tests { }; const LOCAL_ADDR: Ipv4Address = Ipv4Address::new(192, 88, 99, 254); - const LOCAL_PORT: u16 = 8899; + const LOCAL_PORT: u16 = 8900; const PACKETS: TcpPackets = TcpPackets::new(LOCAL_ADDR, LOCAL_PORT); #[tokio::test] async fn accepts_concurrent_connections_with_one_logical_listener() { - let stack = SmolTcpStack::new(LOCAL_ADDR).await.unwrap(); + let stack = SmolTcpStack::new(LOCAL_ADDR, LOCAL_PORT).await.unwrap(); let mut output = stack.take_output_rx().await.unwrap(); let client_addr = Ipv4Address::new(192, 88, 99, 1); diff --git a/easytier/Cargo.toml b/easytier/Cargo.toml index 5f0f25e0..573155eb 100644 --- a/easytier/Cargo.toml +++ b/easytier/Cargo.toml @@ -204,7 +204,7 @@ sys-locale = "0.3" service-manager = { git = "https://github.com/EasyTier/service-manager-rs.git", branch = "main" } -kcp-sys = { git = "https://github.com/EasyTier/kcp-sys", rev = "d7427c22d764deb1860a7d37acc446ed5033464c", optional = true } +kcp-sys = { git = "https://github.com/EasyTier/kcp-sys", rev = "268533568d734ae89dc89603078da3ca522effe1", optional = true } # for dns connector hickory-resolver = { version = "0.25.2", optional = true } diff --git a/easytier/src/tests/three_node.rs b/easytier/src/tests/three_node.rs index 09c50b90..e197c25a 100644 --- a/easytier/src/tests/three_node.rs +++ b/easytier/src/tests/three_node.rs @@ -1454,6 +1454,98 @@ pub async fn subnet_proxy_three_node_test( drop_insts(insts).await; } +#[rstest::rstest] +#[tokio::test] +#[serial_test::serial] +pub async fn subnet_proxy_half_close_test( + #[values("tcp", "kcp", "quic")] transport: &str, + #[values(false, true)] use_smoltcp: bool, +) { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + + let insts = init_three_node_ex( + "udp", + |cfg| { + let mut flags = cfg.get_flags(); + flags.use_smoltcp = use_smoltcp; + if cfg.get_inst_name() == "inst1" { + flags.enable_kcp_proxy = transport == "kcp"; + flags.enable_quic_proxy = transport == "quic"; + } + cfg.set_flags(flags); + if cfg.get_inst_name() == "inst3" { + cfg.add_proxy_cidr( + "10.1.2.0/24".parse().unwrap(), + Some("10.1.3.0/24".parse().unwrap()), + ) + .unwrap(); + } + cfg + }, + false, + ) + .await; + wait_proxy_route_appear( + &insts[0].get_core_instance(), + "10.144.144.3/24", + insts[2].peer_id(), + "10.1.3.0/24", + ) + .await; + + // The advertised route can precede installation of the host TUN route. + wait_for_condition( + || async { ping_test("net_a", "10.1.3.4", None).await }, + Duration::from_secs(5), + ) + .await; + + let result = tokio::time::timeout(Duration::from_secs(10), async { + let listener = NetNS::new(Some("net_d".into())).run(|| { + let listener = std::net::TcpListener::bind("10.1.2.4:22224").unwrap(); + listener.set_nonblocking(true).unwrap(); + tokio::net::TcpListener::from_std(listener).unwrap() + }); + for (source_closes_first, request_size) in + [(true, 0), (false, 0), (true, 64 * 1024), (false, 64 * 1024)] + { + let socket = + NetNS::new(Some("net_a".into())).run(|| tokio::net::TcpSocket::new_v4().unwrap()); + let (client, (server, _)) = tokio::try_join!( + socket.connect("10.1.3.4:22224".parse().unwrap()), + listener.accept(), + ) + .expect("failed to establish the proxied connection"); + let (mut requester, mut responder) = if source_closes_first { + (client, server) + } else { + (server, client) + }; + let request = vec![0x35; request_size]; + let response = vec![0xa7; 128 * 1024]; + tokio::join!( + async { + requester.write_all(&request).await.unwrap(); + requester.shutdown().await.unwrap(); + let mut received = Vec::new(); + requester.read_to_end(&mut received).await.unwrap(); + assert_eq!(received, response); + }, + async { + let mut received = Vec::new(); + responder.read_to_end(&mut received).await.unwrap(); + assert_eq!(received, request); + responder.write_all(&response).await.unwrap(); + responder.shutdown().await.unwrap(); + }, + ); + } + }) + .await; + drop_insts(insts).await; + result.expect("proxy did not forward the response after half-close"); +} + #[rstest::rstest] #[tokio::test] #[serial_test::serial]