Compare commits

...

15 Commits

Author SHA1 Message Date
MOCK
df7947c500 更新tcp日志 2026-03-20 09:04:24 +08:00
nnbcccscdscdsc
4b95d26f13 feat:扩展了结构化日志输出
- 补充了传输层观测指标
  - 增加了 UDP 丢包相关统计
  - README 也同步更新了字段口径
2026-03-17 20:53:43 +08:00
nnbcccscdscdsc
ed1cb20da8 fix:更新.gitignore文件 2026-03-17 19:58:15 +08:00
nnbcccscdscdsc
9f64c9fc49 Stop tracking generated log file 2026-03-17 19:49:05 +08:00
nnbcccscdscdsc
20b2050706 refactor: 收口为hub/peer/bridge 三程序并统一 支持 tcp/udp/kcp" 2026-03-17 16:28:35 +08:00
nnbcccscdscdsc
6c975d9ae3 fix:更新客户端功能 2026-03-16 22:28:05 +08:00
nnbcccscdscdsc
7f2f79e672 feat:实现交叉编译linux x86&arm架构下的client 2026-03-14 08:49:06 +00:00
nnbcccscdscdsc
7ecd8a4ef4 feat: 实现并完成核心功能测试套件
- 编译系统:支持通过 `make clean all` 进行全量编译,生成可执行文件 `omni_client`、`omni_server`、`omni_relay` 和 `omni_test`。
- 客户端-服务端文件传输:支持 TCP/UDP/KCP 协议,已验证文件收发功能(使用 `/tmp/input.bin` 作为测试文件)。
- 服务端指令驱动:服务端可通过控制台发送 ASCII 指令(如 `hello-client`)实时驱动客户端。
- 动态转发功能 (Relay):实现 UDP 协议下的动态目标切换,支持 `show` 查询和 `set` 命令实时修改转发目标(如从 9102 端口切换到 9103 端口)。
- 所有功能已在本地环境(127.0.0.1)通过完整流程验证。
2026-03-13 22:39:41 +08:00
meiqi
4d475f8c92 fix: 重构基础通信框架为异步事件循环,并修复 KCP conv 错位与接收漏斗堵塞问题 2026-03-13 21:00:38 +08:00
meiqi
cf629fa722 fix 2026-03-13 20:03:55 +08:00
meiqi
c819bddbbb log2json 2026-03-13 19:58:29 +08:00
nnbcccscdscdsc
3c1a39c4f4 清理编译中间产物并更新ignore规则 2026-03-13 08:14:02 +00:00
nnbcccscdscdsc
bea7f08f44 Docs:更新ignore文件 2026-03-13 08:07:41 +00:00
meiqi
03e8c7beaa Merge branch 'main' of https://github.com/nnbcccscdscdsc/OmniSocket 2026-03-13 15:53:44 +08:00
nnbcccscdscdsc
3845fc36c0 Docs:更新项目架构 2026-03-13 06:49:21 +00:00
23 changed files with 7397 additions and 390 deletions

24
.gitignore vendored Normal file
View File

@@ -0,0 +1,24 @@
# === 编译产物 (Compiled files) ===
# 忽略所有对象文件
*.o
# === 可执行文件 (Executables) ===
# 建议将所有可执行文件产出到 build/ 目录,然后忽略整个目录
build/
# 如果你直接在根目录生成,手动忽略它们(根据你的项目名修改)
server_pc
client_arm
relay_main
# === 编辑器与 IDE (IDE & Editors) ===
# 忽略 VS Code 配置
.vscode/
# 忽略某些编译数据库
compile_commands.json
.clangd/
# === 日志与测试数据 (Logs & Data) ===
# 忽略程序运行生成的日志文件(csv, log等)
*.log
*.csv
*.jsonl

16
.vscode/launch.json vendored
View File

@@ -1,16 +0,0 @@
{
// 使用 IntelliSense 了解相关属性。
// 悬停以查看现有属性的描述。
// 欲了解更多信息,请访问: https://go.microsoft.com/fwlink/?linkid=830387
"version": "0.2.0",
"configurations": [
{
"name": "OmniSocket",
"type": "lldb",
"request": "launch",
"program": "${workspaceRoot}/build",
"args": [],
"cwd": "${workspaceRoot}"
}
]
}

84
Makefile Normal file
View File

@@ -0,0 +1,84 @@
# 默认本机编译器(可由环境变量覆盖,例如 CC=clang)。
CC ?= gcc
# ARM 交叉编译工具链前缀(可按本地环境替换)。
CROSS_COMPILE ?= arm-linux-gnueabihf-
ARM_CC ?= $(CROSS_COMPILE)gcc
ARM64_CC ?= aarch64-linux-gnu-gcc
# 产物目录(native: build/,arm: build/arm/)。
BUILD_DIR ?= build
# 编译参数:
# - CFLAGS: 优化级别、调试符号、告警、C 标准
# - CPPFLAGS: POSIX 特性宏,确保 clock_gettime 等接口可见
CFLAGS ?= -O2 -g -Wall -Wextra -std=c11
CPPFLAGS ?= -D_DEFAULT_SOURCE -D_POSIX_C_SOURCE=200809L
LDFLAGS ?=
# 链接 pthread。
LDLIBS ?= -lpthread
INCLUDES := -Iinclude
# hub / peer / bridge 共用的核心源码。
PEER_STACK_SRCS := \
src/core/logger.c \
src/core/peer_transport.c \
src/protocols/ikcp.c
# 当前保留的应用入口源文件。
APP_HUB_SRC := src/apps/hub_main.c
APP_PEER_SRC := src/apps/peer_main.c
APP_BRIDGE_SRC := src/apps/bridge_main.c
# 将源文件映射到 BUILD_DIR 下的对象文件路径。
PEER_STACK_OBJS := $(patsubst %.c,$(BUILD_DIR)/%.o,$(PEER_STACK_SRCS))
HUB_OBJ := $(patsubst %.c,$(BUILD_DIR)/%.o,$(APP_HUB_SRC))
PEER_OBJ := $(patsubst %.c,$(BUILD_DIR)/%.o,$(APP_PEER_SRC))
BRIDGE_OBJ := $(patsubst %.c,$(BUILD_DIR)/%.o,$(APP_BRIDGE_SRC))
# 默认构建目标:仅保留 hub / peer / bridge。
TARGETS := \
$(BUILD_DIR)/omni_hub \
$(BUILD_DIR)/omni_peer \
$(BUILD_DIR)/omni_bridge
.PHONY: all arm arm64 clean help
# 本机构建入口。
all: $(TARGETS)
# 各可执行程序链接规则。
$(BUILD_DIR)/omni_hub: $(PEER_STACK_OBJS) $(HUB_OBJ)
$(CC) $(CFLAGS) $(LDFLAGS) -o $@ $^ $(LDLIBS)
$(BUILD_DIR)/omni_peer: $(PEER_STACK_OBJS) $(PEER_OBJ)
$(CC) $(CFLAGS) $(LDFLAGS) -o $@ $^ $(LDLIBS)
$(BUILD_DIR)/omni_bridge: $(PEER_STACK_OBJS) $(BRIDGE_OBJ)
$(CC) $(CFLAGS) $(LDFLAGS) -o $@ $^ $(LDLIBS)
# 通用编译规则:
# - 自动创建对象文件所在目录
# - 编译单个 .c 为 .o
$(BUILD_DIR)/%.o: %.c
@mkdir -p $(dir $@)
$(CC) $(CPPFLAGS) $(CFLAGS) $(INCLUDES) -c $< -o $@
# ARM 构建入口:通过子 make 覆盖 BUILD_DIR 与 CC。
arm:
$(MAKE) BUILD_DIR=build/arm CC=$(ARM_CC) all
# ARM64 构建入口:适合大多数 Jetson 设备。
arm64:
$(MAKE) BUILD_DIR=build/arm64 CC=$(ARM64_CC) all
# 清理构建目录。
clean:
rm -rf build
# 便捷帮助信息。
help:
@echo "make -> build native binaries in build/"
@echo "make arm -> build ARM binaries in build/arm (arm-linux-gnueabihf-gcc)"
@echo "make arm64 -> build ARM64 binaries in build/arm64 (aarch64-linux-gnu-gcc)"
@echo "generated binaries: omni_hub omni_peer omni_bridge"
@echo "make clean -> remove build artifacts"

326
README.md Normal file
View File

@@ -0,0 +1,326 @@
# OmniSocket
OmniSocket 当前包含 3 个核心程序:
- `omni_peer`
- `omni_hub`
- `omni_bridge`
三者统一支持 `tcp | udp | kcp` 三种传输协议。
## 构建
本地构建:
```bash
make
```
生成文件:
- `build/omni_peer`
- `build/omni_hub`
- `build/omni_bridge`
Jetson 场景常用 ARM64 交叉编译:
```bash
make arm64
```
生成文件:
- `build/arm64/omni_peer`
- `build/arm64/omni_hub`
- `build/arm64/omni_bridge`
## 程序说明
### `omni_peer`
`omni_peer` 支持两种工作模式:
- `hub` 模式:连接 `hub` 或 `bridge`
- `direct` 模式:`peer` 之间直接互连
常见用途:
- 发送文本命令
- 发送文件
- 接收文件
### `omni_hub`
`omni_hub` 是中心注册与转发节点,通常部署在公网服务器或中心网络位置。
主要职责:
- 维护 `client_id -> session`
- 转发 `peer` 之间的 `bind / tunnel / status`
### `omni_bridge`
`omni_bridge` 是桥接节点,用于将远端 `peer` 接入上游 `hub`。
主要职责:
- 上游连接 `hub`
- 下游监听本地入口供 `peer` 接入
- 适用于 `A <-> C <-> D <-> B` 这类桥接链路
当前限制:
- 一个 `bridge` 仅支持一个下游 `peer`
- 下游 `peer` 的 `-i` 必须与 `bridge -i` 保持一致
- 上下游必须使用同一种协议,不支持协议转换
- 当前实现更接近“单下游、单身份桥接”,而不是通用多租户中继
## 角色说明
- `A`:本地电脑
- `B`:Jetson
- `C`:公网 Hub 服务器
- `D`:公网 Bridge 服务器
下面示例中的 `<proto>` 可替换为 `tcp`、`udp` 或 `kcp`。
说明:
- 以下 `omni_peer` 示例统一采用长连接交互模式
- 启动后连接会保持,不再使用“传输一次后自动退出”的一次性写法
- 文本和文件传输通过终端中的交互命令完成
- 若启动命令中使用了 `-m` 或 `-F`,程序会执行启动动作模式,而不是进入当前 README 使用的交互模式
## 常用参数与命令
| 分类 | 写法 | 含义 |
| --- | --- | --- |
| 启动参数 | `-i <client_id>` | 当前 `peer` 的逻辑身份。Hub 会根据这个 ID 记录 `client_id -> session` 映射,例如 `-i pc` 表示“我是 pc”,`-i jetson` 表示“我是 jetson”。 |
| 启动参数 | `-b <peer_id>` | 启动后默认绑定的目标 `peer`。例如 `-b jetson` 表示后续直接输入 `send ...` 或 `put ...` 时,默认发给 `jetson`。 |
| 启动参数 | `-d <peer_id>` | 启动动作模式下的显式目标。通常与 `-m` 或 `-F` 配合使用,表示把启动时的那条消息或那个文件直接发给指定目标。 |
| 启动参数 | `-o <output_file>` | 本地接收文件时的落盘路径。收到文件后会写入当前机器上的这个路径,例如 `-o /tmp/from_pc.bin`。 |
| 交互命令 | `bind <peer_id>` | 将默认目标切换到指定 `peer`,后续 `send` 或 `put` 默认发给它。 |
| 交互命令 | `send <text>` | 向当前默认目标发送一条文本消息。 |
| 交互命令 | `say <peer_id> <text>` | 向指定 `peer` 发送一条文本消息,不修改当前默认目标。 |
| 交互命令 | `put <file>` | 将当前机器上的文件发送给当前默认目标。 |
| 交互命令 | `push <peer_id> <file>` | 将当前机器上的文件发送给指定 `peer`,不修改当前默认目标。 |
| 交互命令 | `show` | 显示当前本地状态,例如 `client_id`、当前绑定目标和输出路径。 |
| 交互命令 | `quit` | 退出当前 `omni_peer` 进程。 |
## 场景 1:点对点直传
`A <-> B`
B 端监听:
```bash
./build/omni_peer -M direct -p <proto> -L 9001 -i jetson -o /tmp/from_pc.bin
```
A 端连接:
```bash
./build/omni_peer -M direct -p <proto> -H <B_IP> -P 9001 -i pc -b jetson -o /tmp/from_jetson.bin
```
连接建立后,可在 A 端输入:
```text
send start
put /tmp/input.bin
```
如需反向 `B -> A`,可在 B 端输入:
```text
put /path/to/file.bin
```
## 场景 2:通过 Hub 中转
`A <-> C <-> B`
C 端启动 Hub:
- C 维护着一张 `client_id -> session` 映射表,用于记录谁是 `pc`、谁是 `jetson`,并据此转发 `bind / tunnel / status`
```bash
./build/omni_hub -p <proto> -P 9002
```
B 端连接 Hub:
```bash
./build/omni_peer -p <proto> -H <C_IP> -P 9002 -i jetson -o /tmp/from_pc.bin
```
A 端连接 Hub:
```bash
./build/omni_peer -p <proto> -H <C_IP> -P 9002 -i pc -b jetson -o /tmp/from_jetson.bin
```
连接建立后,可在 A 端输入:
```text
send start
put /tmp/input.bin
```
如需反向 `B -> A`,可在 B 端输入:
```text
bind pc
put /path/to/file.bin
```
## 场景 3:通过 Bridge 桥接
`A <-> C <-> D <-> B`
C 端启动 Hub:
- C 仍然维护 `client_id -> session` 映射表;A 以 `pc` 注册到 C,Bridge 以 `jetson` 这个逻辑身份注册到 C
```bash
./build/omni_hub -p <proto> -P 9003
```
D 端启动 Bridge:
```bash
./build/omni_bridge -p <proto> -H <C_IP> -P 9003 -i jetson -L 9004
```
B 端连接 Bridge:
```bash
./build/omni_peer -p <proto> -H <D_IP> -P 9004 -i jetson -o /tmp/from_pc.bin
```
A 端连接 Hub:
```bash
./build/omni_peer -p <proto> -H <C_IP> -P 9003 -i pc -b jetson -o /tmp/from_jetson.bin
```
连接建立后,可在 A 端输入:
```text
send start
put /tmp/input.bin
```
如需反向 `B -> A`,可在 B 端输入:
```text
bind pc
put /path/to/file.bin
```
说明:
- 该场景已经实现 `A -> C -> D -> B` 与 `B -> D -> C -> A` 的桥接转发
- 但当前 `bridge` 仍是单下游、单身份模型,不是完整的多节点桥接网络
## 日志与指标
### 输出位置
- 终端文本日志:默认输出到 `stderr`,格式为 `key=value`
- 结构化日志:默认追加到当前工作目录下的 `omni_logs.jsonl`
- 性能快照统一写在 `component="perf"` 的 JSONL 记录里
- 终端里会额外出现 `component=perf_udp_loss` 文本行;同一批 UDP 丢包字段已经合并写进对应的 `component="perf"` JSONL 行
- 周期性性能快照大约每 `1s` 打一次;进程退出时会再打一条 `tag="final"`
### 当前文档口径
下面这些说法要以当前实现为准,不要再按旧文档理解:
- `processing / queue / transmission / propagation / end_to_end` 都是“当前实现下的本地观测值或估算值”,不是严格意义上的物理链路精确测量值。
- `processing_*` 当前表示本地应用层处理耗时,主要来自文件分片封装、接收端写盘等路径;它不是“某一块硬件 CPU 的完整开销画像”。
- `queue_*` 和 `transmission_*` 是基于最近活跃窗口的吞吐和缓冲/队列状态反推出来的估算值。最近没有足够流量样本、样本太小或者当前速率太低时,这两个值可能直接为 `0`。
- `propagation_*` 当前来自 `min_rtt_ms / 2` 的估算;如果当前协议没有 RTT 样本,这组字段就是 `0`。
- `end_to_end_*` 当前只在“最终接收文件的 peer”上有值,来源是发送端分片里的 `origin_ts_ms`。发送文件前会先做一次时钟同步,把发送端时间对齐到接收端时钟域;如果同步没建立,这组字段会保持 `0`,而不是给出误导性的跨机结果。
- UDP 丢包统计只在 `UDP 文件接收侧 peer` 上有值;UDP 发送侧、Hub、Bridge 不会产出这组汇总。
- `udp_retrans` 字段当前还没有实现应用层 UDP 重传统计,所以现在始终是 `0`。
- `TCP/KCP` 当前记录的是重传次数/重传字节/累计发送分片,不直接记录“重传频率”这个单独字段。
- `send_buffer_pct_*` / `recv_buffer_pct_*` 是占用率风格的指标,但不保证永远严格落在 `0-100`;尤其 KCP 等待队列超过窗口时,理论上可以大于 `100`。
- `send_call_* / recv_call_* / proto_* / processing_*` 有时会是 `0`,常见原因不是没统计,而是当前时间分辨率是毫秒,很多本地操作小于 `1ms`。
### 基础与吞吐字段
| 文档含义 | JSONL 字段 | 单位 | 当前语义 | 什么时候有值 / 为什么会是 0 |
| --- | --- | --- | --- | --- |
| 时间戳 | `ts_ms` | ms | 这条日志写出的单调时间戳 | 始终有值 |
| 日志分类 | `level` / `component` / `tag` | - | `component="perf"` 表示性能快照,`tag` 常见为 `peer_transport_send`、`peer_transport_recv`、`final` | 始终有值 |
| 身份上下文 | `app` / `proto` / `mode` / `role` / `self_id` | - | 程序名、协议、模式、角色、逻辑 ID | `hub` 的 `self_id` 为空是正常的 |
| 运行时长 | `elapsed_ms` | ms | 当前进程从启动到本次快照的时长 | 始终有值 |
| 累计发送字节 | `bytes_sent` | bytes | 当前进程累计发送的协议帧总字节 | 无发送时为 `0` |
| 累计接收字节 | `bytes_recv` | bytes | 当前进程累计接收的协议帧总字节 | 无接收时为 `0` |
| 发送次数 | `send_count` | count | 当前进程累计发送帧次数 | 无发送时为 `0` |
| 接收次数 | `recv_count` | count | 当前进程累计接收帧次数 | 无接收时为 `0` |
| 瞬时发送带宽 | `tx_current_mbps` | Mbps | 最近一个统计窗口内的发送速率 | 当前窗口没流量时为 `0` |
| 瞬时接收带宽 | `rx_current_mbps` | Mbps | 最近一个统计窗口内的接收速率 | 当前窗口没流量时为 `0` |
| 平均发送带宽 | `tx_avg_mbps` | Mbps | 进程启动到当前的平均发送速率 | 从未发送时为 `0` |
| 平均接收带宽 | `rx_avg_mbps` | Mbps | 进程启动到当前的平均接收速率 | 从未接收时为 `0` |
| 传输进度字节 | `progress_bytes` | bytes | 当前文件传输已完成字节数 | 没有文件传输时为 `0` |
| 总工作量 | `total_work_bytes` | bytes | 当前文件总大小 | 没有文件传输时为 `0` |
| 传输进度百分比 | `progress_pct` | % | `progress_bytes / total_work_bytes * 100` | 没有文件传输时为 `0` |
### 调用耗时与延迟字段
| 文档含义 | JSONL 字段 | 单位 | 当前语义 | 什么时候有值 / 为什么会是 0 |
| --- | --- | --- | --- | --- |
| 应用层发送调用 | `send_call_last_ms` / `send_call_min_ms` / `send_call_max_ms` / `send_call_avg_ms` | ms | `peer_transport_send()` 调用耗时 | 没有发送,或每次发送都小于 `1ms` 时可能为 `0` |
| 应用层接收调用 | `recv_call_last_ms` / `recv_call_min_ms` / `recv_call_max_ms` / `recv_call_avg_ms` | ms | `peer_transport_next_event()` 接收调用耗时 | 没有接收,或每次接收都小于 `1ms` 时可能为 `0` |
| 协议层发送耗时 | `proto_send_avg_ms` | ms | TCP/UDP/KCP 实际 send 路径耗时 EWMA | 发送很快且小于 `1ms` 时常为 `0` |
| 协议层接收耗时 | `proto_recv_avg_ms` | ms | TCP/UDP/KCP 实际 recv 路径耗时 EWMA | 接收很快且小于 `1ms` 时常为 `0` |
| 本地处理耗时 | `processing_avg_ms` / `processing_min_ms` / `processing_max_ms` | ms | 当前进程内的分片封装、写盘等本地处理耗时 | 仅在文件发送/接收路径上采样;操作太快时可能为 `0` |
| 排队延迟估算 | `queue_avg_ms` / `queue_min_ms` / `queue_max_ms` | ms | 根据发送/接收队列字节数和最近活跃窗口速率估算 | 没有足够活跃样本时为 `0`;小样本/低速率时可能偏保守 |
| 传输延迟估算 | `transmission_avg_ms` / `transmission_min_ms` / `transmission_max_ms` | ms | 根据最近活跃窗口速率估算“这些字节推上链路需要多久” | 没有足够活跃样本时为 `0`;小样本/低速率时可能偏保守 |
| 传播延迟估算 | `propagation_avg_ms` / `propagation_min_ms` / `propagation_max_ms` | ms | 基于 `min_rtt_ms / 2` 的估算 | 当前协议没有 RTT 样本时为 `0` |
| 端到端延迟 | `end_to_end_avg_ms` / `end_to_end_min_ms` / `end_to_end_max_ms` | ms | 发送端分片 `origin_ts_ms` 对齐到接收端时钟后,到接收端处理完成时刻的差值 | 只在最终接收文件的 `peer` 上有值;若时钟同步未建立则为 `0` |
### 可靠性字段
| 文档含义 | JSONL 字段 | 单位 | 当前语义 | 什么时候有值 / 为什么会是 0 |
| --- | --- | --- | --- | --- |
| TCP 重传次数 | `tcp_retrans` | count | 内核 `TCP_INFO` 的累计重传次数 | 仅 TCP 有值;未重传或非 TCP 时为 `0` |
| TCP 数据段数 | `tcp_data_segs_out` | count | TCP 累计发送数据段数 | 仅 TCP 有值;非 TCP 时为 `0` |
| TCP 发送字节 | `tcp_data_bytes_sent` | bytes | TCP 累计发送数据字节 | 仅 TCP 有值;非 TCP 时为 `0` |
| TCP 重传字节 | `tcp_retrans_bytes` | bytes | TCP 累计重传字节 | 仅 TCP 有值;未重传或非 TCP 时为 `0` |
| UDP 重传次数 | `udp_retrans` | count | 预留字段,当前未实现 | 当前始终为 `0` |
| KCP 重传次数 | `kcp_retrans` | count | KCP 内部累计重传分片数 | 仅 KCP 有值;未重传或非 KCP 时为 `0` |
| KCP 数据分片数 | `kcp_data_segs_out` | count | KCP 累计发送分片数 | 仅 KCP 有值;非 KCP 时为 `0` |
| KCP 发送字节 | `kcp_data_bytes_sent` | bytes | KCP 累计发送分片字节 | 仅 KCP 有值;非 KCP 时为 `0` |
| KCP 重传字节 | `kcp_retrans_bytes` | bytes | KCP 累计重传字节 | 仅 KCP 有值;未重传或非 KCP 时为 `0` |
| UDP 预期分片数 | `udp_expected_chunks` | count | UDP 文件接收端预期应收到的总分片数 | 仅 UDP 文件接收侧 `peer` 有值;其他角色为 `0` |
| UDP 实收分片数 | `udp_received_chunks` | count | UDP 文件接收端实际收到的唯一分片数 | 仅 UDP 文件接收侧 `peer` 有值;其他角色为 `0` |
| UDP 丢失分片数 | `udp_lost_chunks` | count | 根据 `seq` 推断的缺失分片数 | 完整收到时为 `0`;非 UDP 接收侧也为 `0` |
| UDP 丢包率 | `udp_loss_rate_pct` | % | `udp_lost_chunks / udp_expected_chunks * 100` | 完整收到时为 `0`;非 UDP 接收侧也为 `0` |
| UDP 连续丢包区间数 | `udp_loss_burst_count` | count | 丢包区间个数 | 无丢包时为 `0` |
| UDP 最大连续丢包长度 | `udp_loss_burst_max_len` | count | 单个丢包区间的最大长度 | 无丢包时为 `0` |
| UDP 丢包区间摘要 | `udp_loss_ranges` | csv string | 例如 `4-6,9,12-13` | 无丢包时为空串 |
| UDP 丢失序号样本 | `udp_loss_seq_sample` | csv string | 最多记录一部分缺失 `seq` 样本 | 无丢包时为空串 |
| UDP 接收窗口分布 | `udp_recv_window_dist` | csv string | `window_id:count`,例如 `0:3,1:28` | 仅 UDP 接收侧有值;没有样本时为空串 |
### 资源与算法字段
| 文档含义 | JSONL 字段 | 单位 | 当前语义 | 什么时候有值 / 为什么会是 0 |
| --- | --- | --- | --- | --- |
| 发送缓冲区占用 | `send_buffer_pct_last` / `send_buffer_pct_avg` / `send_buffer_pct_max` | % 风格值 | Socket 或 KCP 发送队列占用率样本 | 没有缓冲采样时为 `0`;KCP 拥塞时可能大于 `100` |
| 接收缓冲区占用 | `recv_buffer_pct_last` / `recv_buffer_pct_avg` / `recv_buffer_pct_max` | % 风格值 | Socket 或 KCP 接收队列占用率样本 | 没有缓冲采样时为 `0` |
| 拥塞窗口 | `cwnd_last` / `cwnd_avg` / `cwnd_max` | 协议窗口大小 | TCP/KCP 当前拥塞窗口样本 | UDP 没有拥塞窗口,因此为 `0` |
| RTT | `last_rtt_ms` / `min_rtt_ms` / `max_rtt_ms` | ms | TCP `TCP_INFO` 或 KCP `rx_srtt` 的 RTT 样本 | UDP 当前没有 RTT 探测,因此为 `0` |
### 实测样本
以下值来自本仓库当前实现的本地回环测试,仅用于说明“字段已经能落到 JSONL 且当前名字是什么”,不是固定性能指标。
| 样本 | 关键字段 | 实测值 |
| --- | --- | --- |
| `TCP` 点对点直传发送端 | `progress_pct` / `cwnd_last` / `last_rtt_ms` / `tcp_data_bytes_sent` | `100` / `10` / `1` / `287984` |
| `UDP` Hub 中转接收端 | `progress_pct` / `udp_expected_chunks` / `udp_received_chunks` / `udp_lost_chunks` / `udp_recv_window_dist` | `100` / `3` / `3` / `0` / `0:3` |
| `KCP` Bridge 桥接节点 | `cwnd_last` / `last_rtt_ms` / `kcp_data_bytes_sent` / `queue_avg_ms` | `2` / `7` / `265` / `33991.342083` |
| `KCP` 最终接收端 | `progress_pct` / `end_to_end_avg_ms` / `cwnd_last` | `100` / `121.333333` / `2` |
### 为什么有些字段“存在但没有值”
| 情况 | 典型字段 | 说明 |
| --- | --- | --- |
| 协议不适用 | `tcp_*` 出现在 UDP/KCP;`kcp_*` 出现在 TCP/UDP;`cwnd_*` 出现在 UDP | 字段统一保留,便于同一份 JSONL 脚本处理;不适用的协议就写 `0` |
| 还没实现 | `udp_retrans` | 现在没有应用层 UDP 重传,所以这个字段只是预留 |
| 当前角色不产出 | `udp_expected_chunks`、`udp_loss_*`、`end_to_end_*` | UDP 丢包统计只在最终接收文件的 UDP peer 上有值;`end_to_end_*` 也主要只在最终接收端有值 |
| 没有采样到 RTT | `last_rtt_ms`、`propagation_*` | UDP 当前没有 RTT 探针;TCP/KCP 只有在拿到对应协议样本后才有值 |
| 没有建立时钟同步 | `end_to_end_*` | 发送端尚未和接收端完成 `TIME_SYNC_*` 探测,或同步结果太旧/未到达 |
| 时间分辨率太粗 | `send_call_*`、`proto_send_avg_ms`、`processing_*` | 当前很多路径按毫秒计时,本地回环下大量操作小于 `1ms`,所以会显示 `0` |
| 当前没有任务 | `progress_*` | 只注册但没有文件传输时,这组字段自然是 `0` |

BIN
ikcp.o

Binary file not shown.

View File

@@ -1,23 +1,209 @@
/* /*
* common.h * common.h
* 全局公共定义:消息头、错误码、通用宏 * 全局公共定义:消息头、错误码、通用宏
*
* 这个头文件承担两类职责:
* 1) 定义跨模块共享的线协议结构,保证 client / server / relay 的编码口径一致。
* 2) 提供与协议无关的基础工具,例如时间戳和 64 位字节序转换。
*/ */
#ifndef OMNISOCKET_COMMON_H #ifndef OMNISOCKET_COMMON_H
#define OMNISOCKET_COMMON_H #define OMNISOCKET_COMMON_H
#include <arpa/inet.h>
#include <stdio.h>
#include <stdint.h> #include <stdint.h>
#include <string.h>
#include <time.h> #include <time.h>
/* 统一的 16 字节消息头(解决 TCP 粘包用) */ /*
* 统一的 16 字节应用层消息头。
* 注意:这不是 TCP/KCP/UDP 的传输层头,而是项目自己定义的“业务帧头”。
*/
typedef struct MsgHeader { typedef struct MsgHeader {
uint32_t magic; /* 固定魔数,用于快速校验 */ uint32_t type; /* 消息类型:文件块/控制指令等 */
uint32_t length; /* 后续负载长度(字节数) */ uint32_t len; /* 后续负载长度(字节数) */
uint64_t seq; /* 序列号或会话内消息 ID */ uint64_t timestamp; /* 发送时间戳(毫秒) */
} MsgHeader; } MsgHeader;
#define MSG_HEADER_SIZE (sizeof(MsgHeader)) /* 16 字节 */ #define MSG_HEADER_SIZE (sizeof(MsgHeader)) /* 16 字节 */
#define MSG_MAGIC 0x4F4D4E49u /* 'OMNI' */
/* 应用层消息类型约定。 */
enum {
MSG_TYPE_FILE_CHUNK = 1, //文件数据分片(TransferChunkMeta + 实际文件字节)
MSG_TYPE_FILE_END = 2, //文件传输结束通知
MSG_TYPE_COMMAND = 3, //控制指令(服务端发起+客户端响应)
MSG_TYPE_TRANSFER_ACK = 4, //传输确认
MSG_TYPE_TIME_SYNC_REQ = 5, //时钟同步探测请求
MSG_TYPE_TIME_SYNC_RESP = 6, //时钟同步探测响应
MSG_TYPE_TIME_SYNC_REPORT = 7, //客户端确认后的最终 offset
MSG_TYPE_PEER_REGISTER = 8, //客户端向 hub 注册自己的逻辑 ID
MSG_TYPE_PEER_BIND = 9, //客户端请求把默认目标绑定到某个 peer_id
MSG_TYPE_PEER_STATUS = 10, //hub 返回给客户端的状态/错误/通知
MSG_TYPE_PEER_TUNNEL = 11, //hub 按 src_id/dst_id 转发的通用隧道消息
MSG_TYPE_RAW = 100 //原始数据(不带业务头,直接透传 payload)
};
/* 默认分片大小(包)上限 */
#define OMNI_DEFAULT_MTU 1400u
/*
* 文件分片元数据:
* - 用于进度统计、端到端延迟估算、UDP 丢包区间分析
* - offset_bytes 让接收端可以按原始偏移写盘,避免 UDP 乱序时直接追加写坏文件
*/
typedef struct TransferChunkMeta {
uint32_t transfer_id; /* 一次文件传输的逻辑 ID,用于区分多次任务。 */
uint32_t seq; /* 当前分片序号,从 1 开始。 */
uint32_t total_chunks; /* 本次传输总分片数,便于接收端判断是否丢片。 */
uint32_t window_id; /* 发送时间窗口编号,用于按秒聚合吞吐/丢包分布。 */
uint64_t total_bytes; /* 原始文件总大小。 */
uint64_t offset_bytes; /* 当前分片在原文件中的字节偏移。 */
uint32_t chunk_bytes; /* 当前分片实际承载的数据字节数。 */
uint32_t reserved; /* 预留字段,当前未使用,发送时写 0。 */
uint64_t origin_ts_ms; /* 分片在发送端最初产生的时间戳,用于端到端延迟估算。 */
} TransferChunkMeta;
#define TRANSFER_CHUNK_META_SIZE (sizeof(TransferChunkMeta))
/*
* 文件发送结束消息。
* 这类消息不携带实际文件内容,而是告诉接收端“数据面已经发完”,
* 便于输出汇总统计并回传 ACK。
*/
typedef struct TransferEndMeta {
uint32_t transfer_id;
uint32_t total_chunks;
uint64_t total_bytes;
uint32_t total_windows;
uint32_t reserved;
} TransferEndMeta;
#define TRANSFER_END_META_SIZE (sizeof(TransferEndMeta))
/*
* 传输结束确认消息。
* 服务端在刷盘和统计完成后回传该结构,客户端据此得知:
* - 哪个 transfer_id 完成了
* - 服务端实际写入了多少字节
* - 这次 ACK 对应的是哪次 FILE_END
*/
typedef struct TransferAckMeta {
uint32_t transfer_id;
uint32_t total_chunks;
uint64_t total_bytes;
uint64_t bytes_written;
uint64_t echoed_end_ts_ms;
} TransferAckMeta;
#define TRANSFER_ACK_META_SIZE (sizeof(TransferAckMeta))
/*
* 时钟同步探测请求。
* 客户端发出探测时带上自己的发送时刻 t0,服务端收到后会立即回包。
*/
typedef struct TimeSyncProbeMeta {
uint32_t probe_id;
uint32_t reserved;
uint64_t client_send_ts_ms;
} TimeSyncProbeMeta;
#define TIME_SYNC_PROBE_META_SIZE (sizeof(TimeSyncProbeMeta))
/*
* 时钟同步响应。
* 服务端把客户端的 t0 原样回显,并补上:
* - t1: 服务端收到请求的本地时间
* - t2: 服务端发出响应的本地时间
*/
typedef struct TimeSyncReplyMeta {
uint32_t probe_id;
uint32_t reserved;
uint64_t client_send_ts_ms;
uint64_t server_recv_ts_ms;
uint64_t server_send_ts_ms;
} TimeSyncReplyMeta;
#define TIME_SYNC_REPLY_META_SIZE (sizeof(TimeSyncReplyMeta))
/*
* 客户端选定最优样本后,把最终 offset 上报给服务端。
* offset 定义为:server_time - client_time。
*/
typedef struct TimeSyncReportMeta {
int64_t server_minus_client_offset_ms;
uint64_t best_rtt_ms;
uint32_t sample_count;
uint32_t reserved;
} TimeSyncReportMeta;
#define TIME_SYNC_REPORT_META_SIZE (sizeof(TimeSyncReportMeta))
#define OMNI_PEER_ID_SIZE 32u
#define OMNI_PEER_STATUS_DETAIL_SIZE 96u
enum {
PEER_STATUS_OK = 0,
PEER_STATUS_ERROR = 1,
PEER_STATUS_REGISTERED = 2,
PEER_STATUS_BOUND = 3,
PEER_STATUS_UNBOUND = 4
};
/*
* peer 注册消息:
* - client_id 是该客户端在 hub 中的逻辑身份
* - 后续 bind / route 都依赖这个 ID,而不是私网 IP
*/
typedef struct PeerRegisterMeta {
char client_id[OMNI_PEER_ID_SIZE];
uint32_t reserved;
} PeerRegisterMeta;
#define PEER_REGISTER_META_SIZE (sizeof(PeerRegisterMeta))
/*
* peer 绑定消息:
* - peer_id 是当前会话的默认目标
* - 绑定成功后,客户端可以直接 send 文本而无需每次重复写目标 ID
*/
typedef struct PeerBindMeta {
char peer_id[OMNI_PEER_ID_SIZE];
uint32_t reserved;
} PeerBindMeta;
#define PEER_BIND_META_SIZE (sizeof(PeerBindMeta))
/*
* hub -> peer 的状态消息:
* - code 区分成功、错误、解绑通知
* - self_id / peer_id 便于客户端更新本地状态
* - detail 留给人读日志和交互终端输出
*/
typedef struct PeerStatusMeta {
uint32_t code;
uint32_t reserved;
char self_id[OMNI_PEER_ID_SIZE];
char peer_id[OMNI_PEER_ID_SIZE];
char detail[OMNI_PEER_STATUS_DETAIL_SIZE];
} PeerStatusMeta;
#define PEER_STATUS_META_SIZE (sizeof(PeerStatusMeta))
/*
* 通用隧道头:
* - src_id 由 hub 在转发时写入真实发送方
* - dst_id 标识目标 peer
* - inner_type 复用现有业务消息类型,后续文件/视频消息也可直接套进来
*/
typedef struct PeerTunnelMeta {
char src_id[OMNI_PEER_ID_SIZE];
char dst_id[OMNI_PEER_ID_SIZE];
uint32_t inner_type;
uint32_t reserved;
} PeerTunnelMeta;
#define PEER_TUNNEL_META_SIZE (sizeof(PeerTunnelMeta))
/* 通用错误码(负数返回表示出错) */ /* 通用错误码(负数返回表示出错) */
enum { enum {
@@ -28,7 +214,7 @@ enum {
OMNI_ERR_TIMEOUT = -4 OMNI_ERR_TIMEOUT = -4
}; };
/* 获取当前单调时间(毫秒),用于延迟统计 */ /* 获取当前单调时间(毫秒),避免系统时间回拨影响延迟统计。 */
static inline uint64_t omni_now_ms(void) static inline uint64_t omni_now_ms(void)
{ {
struct timespec ts; struct timespec ts;
@@ -36,5 +222,342 @@ static inline uint64_t omni_now_ms(void)
return (uint64_t)ts.tv_sec * 1000u + (uint64_t)(ts.tv_nsec / 1000000u); return (uint64_t)ts.tv_sec * 1000u + (uint64_t)(ts.tv_nsec / 1000000u);
} }
#endif /* OMNISOCKET_COMMON_H */ /*
* 手写 64 位字节翻转工具。
* 之所以单独实现,是因为标准库广泛提供 htonl/ntohl,但 64 位版本并非所有平台都统一。
*/
static inline uint64_t omni_bswap64(uint64_t x)
{
return ((x & 0x00000000000000FFull) << 56) |
((x & 0x000000000000FF00ull) << 40) |
((x & 0x0000000000FF0000ull) << 24) |
((x & 0x00000000FF000000ull) << 8) |
((x & 0x000000FF00000000ull) >> 8) |
((x & 0x0000FF0000000000ull) >> 24) |
((x & 0x00FF000000000000ull) >> 40) |
((x & 0xFF00000000000000ull) >> 56);
}
/* 将主机字节序的 64 位整数编码成网络字节序。 */
static inline uint64_t omni_htonll(uint64_t x)
{
#if __BYTE_ORDER__ == __ORDER_LITTLE_ENDIAN__
return omni_bswap64(x);
#else
return x;
#endif
}
/* 64 位网络字节序转回主机字节序;实现上与 htonll 对称。 */
static inline uint64_t omni_ntohll(uint64_t x)
{
return omni_htonll(x);
}
/*
* 保留 int64_t 的原始位模式,便于把 signed offset 按网络字节序放到线协议里。
* 这里用 memcpy,避免依赖实现相关的强制类型转换。
*/
static inline uint64_t omni_i64_bits(int64_t x)
{
uint64_t bits = 0;
memcpy(&bits, &x, sizeof(bits));
return bits;
}
static inline int64_t omni_i64_from_bits(uint64_t bits)
{
int64_t x = 0;
memcpy(&x, &bits, sizeof(x));
return x;
}
static inline void omni_i64_net_encode(int64_t *out_value, int64_t host_value)
{
uint64_t bits = omni_htonll(omni_i64_bits(host_value));
memcpy(out_value, &bits, sizeof(bits));
}
static inline int64_t omni_i64_net_decode(const int64_t *net_value)
{
uint64_t bits = 0;
memcpy(&bits, net_value, sizeof(bits));
return omni_i64_from_bits(omni_ntohll(bits));
}
static inline void omni_copy_fixed_ascii(char *dst, size_t dst_sz, const char *src)
{
if (!dst || dst_sz == 0) {
return;
}
memset(dst, 0, dst_sz);
if (!src) {
return;
}
snprintf(dst, dst_sz, "%s", src);
}
/* 将业务侧的消息头字段编码成网络字节序,便于跨主机传输。 */
static inline void omni_msg_header_encode(MsgHeader *out_hdr,
uint32_t type,
uint32_t len,
uint64_t timestamp_ms)
{
out_hdr->type = htonl(type);
out_hdr->len = htonl(len);
out_hdr->timestamp = omni_htonll(timestamp_ms);
}
/* 将网络帧头解码回本机可直接使用的主机字节序结构。 */
static inline void omni_msg_header_decode(const MsgHeader *net_hdr,
MsgHeader *host_hdr)
{
host_hdr->type = ntohl(net_hdr->type);
host_hdr->len = ntohl(net_hdr->len);
host_hdr->timestamp = omni_ntohll(net_hdr->timestamp);
}
/* 编码单个文件分片的元数据,发送端在每个 chunk 前都要附上它。 */
static inline void omni_transfer_chunk_meta_encode(TransferChunkMeta *out_meta,
uint32_t transfer_id,
uint32_t seq,
uint32_t total_chunks,
uint32_t window_id,
uint64_t total_bytes,
uint64_t offset_bytes,
uint32_t chunk_bytes,
uint64_t origin_ts_ms)
{
out_meta->transfer_id = htonl(transfer_id);
out_meta->seq = htonl(seq);
out_meta->total_chunks = htonl(total_chunks);
out_meta->window_id = htonl(window_id);
out_meta->total_bytes = omni_htonll(total_bytes);
out_meta->offset_bytes = omni_htonll(offset_bytes);
out_meta->chunk_bytes = htonl(chunk_bytes);
out_meta->reserved = 0;
out_meta->origin_ts_ms = omni_htonll(origin_ts_ms);
}
/* 将网络中的分片元数据解码回主机字节序,供接收端统计与写盘使用。 */
static inline void omni_transfer_chunk_meta_decode(const TransferChunkMeta *net_meta,
TransferChunkMeta *host_meta)
{
host_meta->transfer_id = ntohl(net_meta->transfer_id);
host_meta->seq = ntohl(net_meta->seq);
host_meta->total_chunks = ntohl(net_meta->total_chunks);
host_meta->window_id = ntohl(net_meta->window_id);
host_meta->total_bytes = omni_ntohll(net_meta->total_bytes);
host_meta->offset_bytes = omni_ntohll(net_meta->offset_bytes);
host_meta->chunk_bytes = ntohl(net_meta->chunk_bytes);
host_meta->reserved = ntohl(net_meta->reserved);
host_meta->origin_ts_ms = omni_ntohll(net_meta->origin_ts_ms);
}
/* 编码“文件发送完成”消息的元数据。 */
static inline void omni_transfer_end_meta_encode(TransferEndMeta *out_meta,
uint32_t transfer_id,
uint32_t total_chunks,
uint64_t total_bytes,
uint32_t total_windows)
{
out_meta->transfer_id = htonl(transfer_id);
out_meta->total_chunks = htonl(total_chunks);
out_meta->total_bytes = omni_htonll(total_bytes);
out_meta->total_windows = htonl(total_windows);
out_meta->reserved = 0;
}
/* 解码 FILE_END 元数据,接收端据此补全汇总统计。 */
static inline void omni_transfer_end_meta_decode(const TransferEndMeta *net_meta,
TransferEndMeta *host_meta)
{
host_meta->transfer_id = ntohl(net_meta->transfer_id);
host_meta->total_chunks = ntohl(net_meta->total_chunks);
host_meta->total_bytes = omni_ntohll(net_meta->total_bytes);
host_meta->total_windows = ntohl(net_meta->total_windows);
host_meta->reserved = ntohl(net_meta->reserved);
}
/* 编码服务端回给客户端的传输确认信息。 */
static inline void omni_transfer_ack_meta_encode(TransferAckMeta *out_meta,
uint32_t transfer_id,
uint32_t total_chunks,
uint64_t total_bytes,
uint64_t bytes_written,
uint64_t echoed_end_ts_ms)
{
out_meta->transfer_id = htonl(transfer_id);
out_meta->total_chunks = htonl(total_chunks);
out_meta->total_bytes = omni_htonll(total_bytes);
out_meta->bytes_written = omni_htonll(bytes_written);
out_meta->echoed_end_ts_ms = omni_htonll(echoed_end_ts_ms);
}
/* 解码 ACK 元数据,客户端用它计算确认 RTT 和最终落盘结果。 */
static inline void omni_transfer_ack_meta_decode(const TransferAckMeta *net_meta,
TransferAckMeta *host_meta)
{
host_meta->transfer_id = ntohl(net_meta->transfer_id);
host_meta->total_chunks = ntohl(net_meta->total_chunks);
host_meta->total_bytes = omni_ntohll(net_meta->total_bytes);
host_meta->bytes_written = omni_ntohll(net_meta->bytes_written);
host_meta->echoed_end_ts_ms = omni_ntohll(net_meta->echoed_end_ts_ms);
}
/* 编码时钟同步探测请求。 */
static inline void omni_time_sync_probe_meta_encode(TimeSyncProbeMeta *out_meta,
uint32_t probe_id,
uint64_t client_send_ts_ms)
{
out_meta->probe_id = htonl(probe_id);
out_meta->reserved = 0;
out_meta->client_send_ts_ms = omni_htonll(client_send_ts_ms);
}
/* 解码时钟同步探测请求。 */
static inline void omni_time_sync_probe_meta_decode(const TimeSyncProbeMeta *net_meta,
TimeSyncProbeMeta *host_meta)
{
host_meta->probe_id = ntohl(net_meta->probe_id);
host_meta->reserved = ntohl(net_meta->reserved);
host_meta->client_send_ts_ms = omni_ntohll(net_meta->client_send_ts_ms);
}
/* 编码时钟同步响应。 */
static inline void omni_time_sync_reply_meta_encode(TimeSyncReplyMeta *out_meta,
uint32_t probe_id,
uint64_t client_send_ts_ms,
uint64_t server_recv_ts_ms,
uint64_t server_send_ts_ms)
{
out_meta->probe_id = htonl(probe_id);
out_meta->reserved = 0;
out_meta->client_send_ts_ms = omni_htonll(client_send_ts_ms);
out_meta->server_recv_ts_ms = omni_htonll(server_recv_ts_ms);
out_meta->server_send_ts_ms = omni_htonll(server_send_ts_ms);
}
/* 解码时钟同步响应。 */
static inline void omni_time_sync_reply_meta_decode(const TimeSyncReplyMeta *net_meta,
TimeSyncReplyMeta *host_meta)
{
host_meta->probe_id = ntohl(net_meta->probe_id);
host_meta->reserved = ntohl(net_meta->reserved);
host_meta->client_send_ts_ms = omni_ntohll(net_meta->client_send_ts_ms);
host_meta->server_recv_ts_ms = omni_ntohll(net_meta->server_recv_ts_ms);
host_meta->server_send_ts_ms = omni_ntohll(net_meta->server_send_ts_ms);
}
/* 编码客户端最终确认的 offset。 */
static inline void omni_time_sync_report_meta_encode(TimeSyncReportMeta *out_meta,
int64_t server_minus_client_offset_ms,
uint64_t best_rtt_ms,
uint32_t sample_count)
{
omni_i64_net_encode(&out_meta->server_minus_client_offset_ms,
server_minus_client_offset_ms);
out_meta->best_rtt_ms = omni_htonll(best_rtt_ms);
out_meta->sample_count = htonl(sample_count);
out_meta->reserved = 0;
}
/* 解码客户端最终确认的 offset。 */
static inline void omni_time_sync_report_meta_decode(const TimeSyncReportMeta *net_meta,
TimeSyncReportMeta *host_meta)
{
host_meta->server_minus_client_offset_ms =
omni_i64_net_decode(&net_meta->server_minus_client_offset_ms);
host_meta->best_rtt_ms = omni_ntohll(net_meta->best_rtt_ms);
host_meta->sample_count = ntohl(net_meta->sample_count);
host_meta->reserved = ntohl(net_meta->reserved);
}
/* 编码 peer 注册消息。 */
static inline void omni_peer_register_meta_encode(PeerRegisterMeta *out_meta,
const char *client_id)
{
memset(out_meta, 0, sizeof(*out_meta));
omni_copy_fixed_ascii(out_meta->client_id, sizeof(out_meta->client_id), client_id);
out_meta->reserved = 0;
}
/* 解码 peer 注册消息。 */
static inline void omni_peer_register_meta_decode(const PeerRegisterMeta *net_meta,
PeerRegisterMeta *host_meta)
{
memset(host_meta, 0, sizeof(*host_meta));
omni_copy_fixed_ascii(host_meta->client_id, sizeof(host_meta->client_id), net_meta->client_id);
host_meta->reserved = ntohl(net_meta->reserved);
}
/* 编码 peer bind 消息。 */
static inline void omni_peer_bind_meta_encode(PeerBindMeta *out_meta,
const char *peer_id)
{
memset(out_meta, 0, sizeof(*out_meta));
omni_copy_fixed_ascii(out_meta->peer_id, sizeof(out_meta->peer_id), peer_id);
out_meta->reserved = 0;
}
/* 解码 peer bind 消息。 */
static inline void omni_peer_bind_meta_decode(const PeerBindMeta *net_meta,
PeerBindMeta *host_meta)
{
memset(host_meta, 0, sizeof(*host_meta));
omni_copy_fixed_ascii(host_meta->peer_id, sizeof(host_meta->peer_id), net_meta->peer_id);
host_meta->reserved = ntohl(net_meta->reserved);
}
/* 编码 hub 状态消息。 */
static inline void omni_peer_status_meta_encode(PeerStatusMeta *out_meta,
uint32_t code,
const char *self_id,
const char *peer_id,
const char *detail)
{
memset(out_meta, 0, sizeof(*out_meta));
out_meta->code = htonl(code);
out_meta->reserved = 0;
omni_copy_fixed_ascii(out_meta->self_id, sizeof(out_meta->self_id), self_id);
omni_copy_fixed_ascii(out_meta->peer_id, sizeof(out_meta->peer_id), peer_id);
omni_copy_fixed_ascii(out_meta->detail, sizeof(out_meta->detail), detail);
}
/* 解码 hub 状态消息。 */
static inline void omni_peer_status_meta_decode(const PeerStatusMeta *net_meta,
PeerStatusMeta *host_meta)
{
memset(host_meta, 0, sizeof(*host_meta));
host_meta->code = ntohl(net_meta->code);
host_meta->reserved = ntohl(net_meta->reserved);
omni_copy_fixed_ascii(host_meta->self_id, sizeof(host_meta->self_id), net_meta->self_id);
omni_copy_fixed_ascii(host_meta->peer_id, sizeof(host_meta->peer_id), net_meta->peer_id);
omni_copy_fixed_ascii(host_meta->detail, sizeof(host_meta->detail), net_meta->detail);
}
/* 编码 peer tunnel 头。 */
static inline void omni_peer_tunnel_meta_encode(PeerTunnelMeta *out_meta,
const char *src_id,
const char *dst_id,
uint32_t inner_type)
{
memset(out_meta, 0, sizeof(*out_meta));
omni_copy_fixed_ascii(out_meta->src_id, sizeof(out_meta->src_id), src_id);
omni_copy_fixed_ascii(out_meta->dst_id, sizeof(out_meta->dst_id), dst_id);
out_meta->inner_type = htonl(inner_type);
out_meta->reserved = 0;
}
/* 解码 peer tunnel 头。 */
static inline void omni_peer_tunnel_meta_decode(const PeerTunnelMeta *net_meta,
PeerTunnelMeta *host_meta)
{
memset(host_meta, 0, sizeof(*host_meta));
omni_copy_fixed_ascii(host_meta->src_id, sizeof(host_meta->src_id), net_meta->src_id);
omni_copy_fixed_ascii(host_meta->dst_id, sizeof(host_meta->dst_id), net_meta->dst_id);
host_meta->inner_type = ntohl(net_meta->inner_type);
host_meta->reserved = ntohl(net_meta->reserved);
}
#endif /* OMNISOCKET_COMMON_H */

View File

@@ -9,28 +9,103 @@
#include <stddef.h> #include <stddef.h>
#include <stdint.h> #include <stdint.h>
typedef struct OmniMetricSummary {
uint64_t count;
double last;
double min;
double max;
double sum;
} OmniMetricSummary;
#define OMNI_LOGGER_UDP_RANGES_SIZE 512u
#define OMNI_LOGGER_UDP_SEQ_SAMPLE_SIZE 512u
#define OMNI_LOGGER_UDP_WINDOW_DIST_SIZE 512u
/* 通过该结构体收集全局统计信息 */ /* 通过该结构体收集全局统计信息 */
typedef struct OmniStats { typedef struct OmniStats {
uint64_t start_ms; /* 起始时间(毫秒) */ uint64_t start_ms; /* 起始时间(毫秒) */
uint64_t last_report_ms; /* 上一次打印日志时间 */ uint64_t last_report_ms; /* 上一次打印日志时间 */
uint64_t window_start_ms; /* 当前吞吐统计窗口起点 */
uint64_t bytes_sent; /* 发送总字节数 */ uint64_t bytes_sent; /* 发送总字节数 */
uint64_t bytes_recv; /* 接收总字节数 */ uint64_t bytes_recv; /* 接收总字节数 */
uint64_t window_bytes_sent; /* 当前 1 秒窗口发送字节数 */
uint64_t window_bytes_recv; /* 当前 1 秒窗口接收字节数 */
uint64_t delay_window_start_send_ms; /* 延时估算使用的最近发送活跃窗口起点 */
uint64_t delay_window_start_recv_ms; /* 延时估算使用的最近接收活跃窗口起点 */
uint64_t delay_window_bytes_sent; /* 延时估算使用的最近发送活跃窗口字节数 */
uint64_t delay_window_bytes_recv; /* 延时估算使用的最近接收活跃窗口字节数 */
uint64_t last_send_activity_ms; /* 最近一次发送活跃时间 */
uint64_t last_recv_activity_ms; /* 最近一次接收活跃时间 */
uint64_t send_count; /* 调用 omni_send 次数 */ uint64_t send_count; /* 调用 omni_send 次数 */
uint64_t recv_count; /* 调用 omni_recv 次数 */ uint64_t recv_count; /* 调用 omni_recv 次数 */
uint64_t last_rtt_ms; /* 最近一次 RTT */ uint64_t last_rtt_ms; /* 最近一次 RTT */
uint64_t min_rtt_ms; /* 最小 RTT(更接近链路基线) */
uint64_t max_rtt_ms; /* 最大 RTT */ uint64_t max_rtt_ms; /* 最大 RTT */
uint64_t tcp_retrans; /* 预留:TCP 重传统计(如可从内核获取) */ uint64_t tcp_retrans; /* 预留:TCP 重传统计(如可从内核获取) */
uint64_t udp_retrans; /* UDP 上层重传次数 */ uint64_t udp_retrans; /* UDP 上层重传次数 */
uint64_t kcp_retrans; /* KCP 内部重传次数(可从 ikcp 统计) */ uint64_t kcp_retrans; /* KCP 内部重传次数(可从 ikcp 统计) */
uint64_t tcp_data_segs_out; /* TCP 累计发送的数据段数(含重传) */
uint64_t tcp_data_bytes_sent; /* TCP 累计发送的数据字节(含重传) */
uint64_t tcp_retrans_bytes; /* TCP 累计重传的数据字节 */
uint64_t kcp_data_segs_out; /* KCP 累计发送的数据分片数(含重传) */
uint64_t kcp_data_bytes_sent; /* KCP 累计发送的数据字节(含重传) */
uint64_t kcp_retrans_bytes; /* KCP 累计重传的数据字节 */
uint64_t udp_expected_chunks; /* UDP 文件接收侧预期分片数 */
uint64_t udp_received_chunks; /* UDP 文件接收侧实际收到的分片数 */
uint64_t udp_lost_chunks; /* UDP 文件接收侧推断丢失的分片数 */
uint64_t udp_loss_burst_count; /* UDP 丢包区间数量 */
uint64_t udp_loss_burst_max_len; /* UDP 最大连续丢包长度 */
double udp_loss_rate_pct; /* UDP 丢包率(百分比) */
char udp_loss_ranges[OMNI_LOGGER_UDP_RANGES_SIZE]; /* UDP 丢包区间摘要 */
char udp_loss_seq_sample[OMNI_LOGGER_UDP_SEQ_SAMPLE_SIZE]; /* UDP 丢包序号样本 */
char udp_recv_window_dist[OMNI_LOGGER_UDP_WINDOW_DIST_SIZE]; /* UDP 接收窗口分布 */
/* 延迟/耗时统计(单位:毫秒) */
double send_call_avg_ms; /* omni_send 平均耗时(EWMA) */
double recv_call_avg_ms; /* omni_recv 平均耗时(EWMA) */
double proto_send_avg_ms; /* 协议 send() 平均耗时(EWMA) */
double proto_recv_avg_ms; /* 协议 recv() 平均耗时(EWMA) */
uint64_t send_call_min_ms;
uint64_t send_call_max_ms;
uint64_t recv_call_min_ms;
uint64_t recv_call_max_ms;
uint64_t last_send_call_ms;
uint64_t last_recv_call_ms;
uint64_t total_work_bytes; /* 本次任务总大小(文件总大小等) */
uint64_t progress_bytes; /* 当前进度字节数 */
double tx_current_mbps; /* 最近 1 秒发送速率(Mbps) */
double rx_current_mbps; /* 最近 1 秒接收速率(Mbps) */
double tx_avg_mbps; /* 从开始到当前的平均发送速率(Mbps) */
double rx_avg_mbps; /* 从开始到当前的平均接收速率(Mbps) */
OmniMetricSummary processing_delay_ms; // 上层处理耗时(如文件读写、加解密等)
OmniMetricSummary queue_delay_ms; // 排队延迟
OmniMetricSummary transmission_delay_ms; // 传输延迟
OmniMetricSummary propagation_delay_ms; // 传播延迟
OmniMetricSummary end_to_end_delay_ms; // 端到端延迟
OmniMetricSummary send_buffer_pct; // 发送缓冲区占用率
OmniMetricSummary recv_buffer_pct; // 接收缓冲区占用率
OmniMetricSummary cwnd; // 拥塞窗口大小
} OmniStats; } OmniStats;
/* 初始化统计模块,在程序启动时调用一次 */ /* 初始化统计模块,在程序启动时调用一次 */
void logger_init(void); void logger_init(void);
/* 设置当前进程的日志上下文,便于 perf/jsonl 日志携带协议与节点信息。 */
void logger_set_context(const char *app,
const char *proto,
const char *mode,
const char *role,
const char *self_id);
/* 记录一次发送/接收 */ /* 记录一次发送/接收 */
void logger_on_send(size_t bytes); void logger_on_send(size_t bytes);
void logger_on_recv(size_t bytes); void logger_on_recv(size_t bytes);
@@ -38,13 +113,55 @@ void logger_on_recv(size_t bytes);
/* 记录一次 RTT(由上层在合适时机调用) */ /* 记录一次 RTT(由上层在合适时机调用) */
void logger_on_rtt(uint64_t rtt_ms); void logger_on_rtt(uint64_t rtt_ms);
/* 记录 KCP 重传次数变化(可在 KCP 更新循环中调用) */ /* 记录 TCP 传输层累计快照(通常来自 TCP_INFO)。 */
void logger_on_kcp_retrans(uint64_t delta); void logger_on_tcp_transport(uint64_t total_retrans,
uint64_t data_segs_out,
uint64_t data_bytes_sent,
uint64_t retrans_bytes);
/* 记录 KCP 数据分片首次发送。 */
void logger_on_kcp_tx(uint64_t segs, uint64_t bytes);
/* 记录 KCP 数据分片重传。 */
void logger_on_kcp_retrans(uint64_t segs, uint64_t bytes);
/* 记录一次耗时(ms) */
void logger_on_send_call_latency(uint64_t ms);
void logger_on_recv_call_latency(uint64_t ms);
void logger_on_proto_send_latency(uint64_t ms);
void logger_on_proto_recv_latency(uint64_t ms);
void logger_on_processing_latency(double ms);
void logger_on_queue_delay_est(double ms);
void logger_on_transmission_delay_est(double ms);
void logger_on_propagation_delay_est(double ms);
void logger_on_end_to_end_latency(double ms);
void logger_on_send_queue_bytes(size_t bytes);
void logger_on_recv_queue_bytes(size_t bytes);
void logger_on_send_transmission_bytes(size_t bytes);
void logger_on_recv_transmission_bytes(size_t bytes);
void logger_on_buffer_status(double send_pct, double recv_pct);
void logger_on_cwnd(double cwnd);
/* 记录任务总量与当前进度。 */
void logger_set_transfer_total(uint64_t total_bytes);
void logger_set_progress(uint64_t progress_bytes);
void logger_reset_transfer_observability(void);
void logger_on_udp_loss_summary(uint64_t expected_chunks,
uint64_t received_chunks,
uint64_t lost_chunks,
uint64_t burst_count,
uint64_t burst_max_len,
const char *ranges,
const char *seq_sample,
const char *recv_window_dist);
/* 计算当前吞吐量(返回:字节/秒) */ /* 计算当前吞吐量(返回:字节/秒) */
double logger_calculate_throughput(void); double logger_calculate_throughput(void);
/* 打印一条结构化性能日志(例如每隔若干秒调用) */ /* 在 1 秒窗口到期时打印周期性性能日志。 */
void logger_maybe_print_performance_log(const char *tag);
/* 强制打印一条结构化性能日志(例如最终汇总前调用) */
void logger_print_performance_log(const char *tag); void logger_print_performance_log(const char *tag);
/* 结构化通用日志(key=value 形式) */ /* 结构化通用日志(key=value 形式) */
@@ -55,4 +172,3 @@ void logger_log(const char *level, const char *component,
OmniStats logger_get_snapshot(void); OmniStats logger_get_snapshot(void);
#endif /* OMNISOCKET_LOGGER_H */ #endif /* OMNISOCKET_LOGGER_H */

59
include/peer_transport.h Normal file
View File

@@ -0,0 +1,59 @@
#ifndef OMNISOCKET_PEER_TRANSPORT_H
#define OMNISOCKET_PEER_TRANSPORT_H
#include "common.h"
#include "network.h"
#include <stddef.h>
#include <stdint.h>
typedef struct PeerTransport PeerTransport;
typedef struct PeerTransportSession PeerTransportSession;
typedef struct PeerTransportEvent {
int kind;
PeerTransportSession *session;
MsgHeader header;
} PeerTransportEvent;
enum {
PEER_TRANSPORT_EVENT_NONE = 0,
PEER_TRANSPORT_EVENT_MESSAGE = 1,
PEER_TRANSPORT_EVENT_CLOSED = 2
};
PeerTransport *peer_transport_open(OmniRole role,
OmniProtocol proto,
const char *bind_ip,
uint16_t bind_port,
const char *peer_ip,
uint16_t peer_port);
void peer_transport_close(PeerTransport *transport);
PeerTransportSession *peer_transport_default_session(PeerTransport *transport);
int peer_transport_send(PeerTransport *transport,
PeerTransportSession *session,
uint32_t type,
const void *payload,
uint32_t payload_len);
int peer_transport_next_event(PeerTransport *transport,
PeerTransportEvent *event,
uint8_t *payload_buf,
size_t payload_cap,
int timeout_ms);
void peer_transport_close_session(PeerTransport *transport,
PeerTransportSession *session);
const char *peer_transport_session_remote(const PeerTransportSession *session,
char *buf,
size_t buf_sz);
uint16_t peer_transport_session_remote_port(const PeerTransportSession *session);
const char *peer_transport_proto_name(OmniProtocol proto);
#endif

Binary file not shown.

BIN
logger.o

Binary file not shown.

BIN
network.o

Binary file not shown.

View File

@@ -0,0 +1,49 @@
#!/usr/bin/env bash
# 本机 hub/peer smoke 测试:
# - 启动 1 个 omni_hub
# - 启动 2 个 omni_peer(beta 常驻,alpha 发一条消息给 beta)
# - 校验 beta 日志里确实收到了 alpha 的消息
set -euo pipefail
ROOT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
BUILD_DIR="$ROOT_DIR/build"
TMP_DIR="$(mktemp -d /tmp/omnisocket-peer-smoke.XXXXXX)"
PORT=$((30000 + (RANDOM % 20000)))
PIDS=()
cleanup() {
for pid in "${PIDS[@]:-}"; do
kill "$pid" 2>/dev/null || true
wait "$pid" 2>/dev/null || true
done
rm -rf "$TMP_DIR"
}
trap cleanup EXIT
log() {
printf '[peer-smoke] %s\n' "$1"
}
log "building hub/peer binaries"
make -C "$ROOT_DIR" build/omni_hub build/omni_peer >/dev/null
log "starting hub on port $PORT"
"$BUILD_DIR/omni_hub" -P "$PORT" >"$TMP_DIR/hub.log" 2>&1 &
HUB_PID=$!
PIDS+=("$HUB_PID")
sleep 1
log "starting beta peer"
"$BUILD_DIR/omni_peer" -H 127.0.0.1 -P "$PORT" -i beta -w 4 >"$TMP_DIR/beta.log" 2>&1 &
BETA_PID=$!
PIDS+=("$BETA_PID")
sleep 1
log "starting alpha peer and sending command"
"$BUILD_DIR/omni_peer" -H 127.0.0.1 -P "$PORT" -i alpha -b beta -d beta -m "hello-from-alpha" -w 2 >"$TMP_DIR/alpha.log" 2>&1
wait "$BETA_PID"
grep -q "\[peer alpha -> beta\] hello-from-alpha" "$TMP_DIR/beta.log"
log "peer command tunnel passed"

411
src/apps/bridge_main.c Normal file
View File

@@ -0,0 +1,411 @@
/*
* bridge_main.c
* 固定多跳桥接:
* - 上游作为一个 peer 主动连接远端 hub
* - 下游作为一个轻量 hub 接入本地 peer
* - 将 bind / tunnel / status 在上下游之间转发
*/
#include "common.h"
#include "logger.h"
#include "peer_transport.h"
#include <signal.h>
#include <stdint.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#define BRIDGE_MAX_PAYLOAD (PEER_TUNNEL_META_SIZE + 65536u)
typedef struct BridgeRuntime {
PeerTransport *upstream;
PeerTransport *downstream;
PeerTransportSession *downstream_session;
char client_id[OMNI_PEER_ID_SIZE];
} BridgeRuntime;
static volatile sig_atomic_t g_stop = 0;
static void on_signal(int signo)
{
(void)signo;
g_stop = 1;
}
static void install_signal_handlers(void)
{
struct sigaction sa;
memset(&sa, 0, sizeof(sa));
sa.sa_handler = on_signal;
sigemptyset(&sa.sa_mask);
(void)sigaction(SIGINT, &sa, NULL);
(void)sigaction(SIGTERM, &sa, NULL);
(void)signal(SIGPIPE, SIG_IGN);
}
static void usage(const char *prog)
{
fprintf(stderr,
"Usage:\n"
" %s -H <upstream_hub_ip> -P <upstream_hub_port> -i <client_id> -L <listen_port>\n"
" [-b <bind_ip>] [-p tcp|udp|kcp]\n",
prog);
}
static int parse_proto(const char *s, OmniProtocol *out_proto)
{
if (!s || !out_proto) {
return 0;
}
if (strcmp(s, "tcp") == 0) {
*out_proto = OMNI_PROTO_TCP;
return 1;
}
if (strcmp(s, "udp") == 0) {
*out_proto = OMNI_PROTO_UDP;
return 1;
}
if (strcmp(s, "kcp") == 0) {
*out_proto = OMNI_PROTO_KCP;
return 1;
}
return 0;
}
static int peer_id_is_valid(const char *id)
{
size_t len = 0;
if (!id || !id[0]) {
return 0;
}
for (len = 0; id[len] != '\0'; ++len) {
unsigned char ch = (unsigned char)id[len];
if (len + 1u >= OMNI_PEER_ID_SIZE) {
return 0;
}
if (!((ch >= 'a' && ch <= 'z') ||
(ch >= 'A' && ch <= 'Z') ||
(ch >= '0' && ch <= '9') ||
ch == '_' || ch == '-' || ch == '.')) {
return 0;
}
}
return 1;
}
static int send_status_to_downstream(BridgeRuntime *rt,
uint32_t code,
const char *self_id,
const char *peer_id,
const char *detail)
{
PeerStatusMeta status_meta;
if (!rt || !rt->downstream || !rt->downstream_session) {
return OMNI_ERR_PARAM;
}
omni_peer_status_meta_encode(&status_meta, code, self_id, peer_id, detail);
return peer_transport_send(rt->downstream,
rt->downstream_session,
MSG_TYPE_PEER_STATUS,
&status_meta,
PEER_STATUS_META_SIZE);
}
static int send_upstream_register(BridgeRuntime *rt)
{
PeerRegisterMeta meta;
if (!rt || !rt->upstream) {
return OMNI_ERR_PARAM;
}
omni_peer_register_meta_encode(&meta, rt->client_id);
return peer_transport_send(rt->upstream,
NULL,
MSG_TYPE_PEER_REGISTER,
&meta,
PEER_REGISTER_META_SIZE);
}
static int handle_upstream_event(BridgeRuntime *rt,
const PeerTransportEvent *event,
const uint8_t *payload)
{
if (!rt || !event) {
return OMNI_ERR_PARAM;
}
if (event->kind == PEER_TRANSPORT_EVENT_CLOSED) {
logger_log("INFO", "bridge", "upstream_closed");
return OMNI_ERR_IO;
}
switch (event->header.type) {
case MSG_TYPE_PEER_STATUS:
case MSG_TYPE_PEER_TUNNEL:
if (!rt->downstream_session) {
logger_log("WARN", "bridge",
"downstream_forward_skipped type=%u len=%u",
(unsigned)event->header.type,
(unsigned)event->header.len);
return OMNI_OK;
}
return peer_transport_send(rt->downstream,
rt->downstream_session,
event->header.type,
payload,
event->header.len);
default:
logger_log("WARN", "bridge",
"unexpected_upstream_type=%u len=%u",
(unsigned)event->header.type,
(unsigned)event->header.len);
return OMNI_OK;
}
}
static int handle_downstream_register(BridgeRuntime *rt,
PeerTransportSession *session,
const uint8_t *payload,
uint32_t payload_len)
{
PeerRegisterMeta register_meta;
if (!rt || !session) {
return OMNI_ERR_PARAM;
}
if (payload_len < PEER_REGISTER_META_SIZE) {
rt->downstream_session = session;
return send_status_to_downstream(rt,
PEER_STATUS_ERROR,
NULL,
NULL,
"short_register_payload");
}
omni_peer_register_meta_decode((const PeerRegisterMeta *)payload, &register_meta);
if (strcmp(register_meta.client_id, rt->client_id) != 0) {
rt->downstream_session = session;
logger_log("WARN", "bridge",
"downstream_register_mismatch got=%s expect=%s",
register_meta.client_id,
rt->client_id);
return send_status_to_downstream(rt,
PEER_STATUS_ERROR,
register_meta.client_id,
NULL,
"client_id_mismatch");
}
rt->downstream_session = session;
return send_status_to_downstream(rt,
PEER_STATUS_REGISTERED,
rt->client_id,
NULL,
"bridge_register_ok");
}
static int handle_downstream_event(BridgeRuntime *rt,
const PeerTransportEvent *event,
const uint8_t *payload)
{
if (!rt || !event) {
return OMNI_ERR_PARAM;
}
if (event->kind == PEER_TRANSPORT_EVENT_CLOSED) {
if (rt->downstream_session == event->session) {
logger_log("INFO", "bridge", "downstream_closed");
rt->downstream_session = NULL;
}
peer_transport_close_session(rt->downstream, event->session);
return OMNI_OK;
}
if (rt->downstream_session && rt->downstream_session != event->session) {
PeerStatusMeta status_meta;
omni_peer_status_meta_encode(&status_meta,
PEER_STATUS_ERROR,
rt->client_id,
NULL,
"bridge_busy");
(void)peer_transport_send(rt->downstream,
event->session,
MSG_TYPE_PEER_STATUS,
&status_meta,
PEER_STATUS_META_SIZE);
peer_transport_close_session(rt->downstream, event->session);
logger_log("WARN", "bridge", "reject_extra_downstream");
return OMNI_OK;
}
switch (event->header.type) {
case MSG_TYPE_PEER_REGISTER:
return handle_downstream_register(rt,
event->session,
payload,
event->header.len);
case MSG_TYPE_PEER_BIND:
case MSG_TYPE_PEER_TUNNEL:
if (rt->downstream_session != event->session) {
rt->downstream_session = event->session;
(void)send_status_to_downstream(rt,
PEER_STATUS_ERROR,
rt->client_id,
NULL,
"register_first");
return OMNI_OK;
}
if (peer_transport_send(rt->upstream,
NULL,
event->header.type,
payload,
event->header.len) != OMNI_OK) {
(void)send_status_to_downstream(rt,
PEER_STATUS_ERROR,
rt->client_id,
NULL,
"upstream_forward_failed");
}
return OMNI_OK;
default:
logger_log("WARN", "bridge",
"unexpected_downstream_type=%u len=%u",
(unsigned)event->header.type,
(unsigned)event->header.len);
return OMNI_OK;
}
}
int main(int argc, char **argv)
{
const char *upstream_ip = NULL;
const char *bind_ip = NULL;
const char *client_id = NULL;
const char *proto_str = "tcp";
OmniProtocol proto = OMNI_PROTO_TCP;
int upstream_port = 0;
int listen_port = 0;
int opt;
BridgeRuntime rt;
uint8_t upstream_payload[BRIDGE_MAX_PAYLOAD];
uint8_t downstream_payload[BRIDGE_MAX_PAYLOAD];
while ((opt = getopt(argc, argv, "H:P:i:L:b:p:")) != -1) {
switch (opt) {
case 'H':
upstream_ip = optarg;
break;
case 'P':
upstream_port = atoi(optarg);
break;
case 'i':
client_id = optarg;
break;
case 'L':
listen_port = atoi(optarg);
break;
case 'b':
bind_ip = optarg;
break;
case 'p':
proto_str = optarg;
break;
default:
usage(argv[0]);
return 1;
}
}
if (!upstream_ip || upstream_port <= 0 || listen_port <= 0 ||
!peer_id_is_valid(client_id) || !parse_proto(proto_str, &proto)) {
usage(argv[0]);
return 1;
}
logger_init();
install_signal_handlers();
logger_set_context("bridge", proto_str, "bridge", "relay", client_id);
memset(&rt, 0, sizeof(rt));
omni_copy_fixed_ascii(rt.client_id, sizeof(rt.client_id), client_id);
rt.upstream = peer_transport_open(OMNI_ROLE_CLIENT,
proto,
NULL,
0,
upstream_ip,
(uint16_t)upstream_port);
if (!rt.upstream) {
perror("bridge upstream");
return 1;
}
if (send_upstream_register(&rt) != OMNI_OK) {
fprintf(stderr, "bridge upstream register failed\n");
peer_transport_close(rt.upstream);
return 1;
}
rt.downstream = peer_transport_open(OMNI_ROLE_SERVER,
proto,
bind_ip,
(uint16_t)listen_port,
NULL,
0);
if (!rt.downstream) {
perror("bridge downstream");
peer_transport_close(rt.upstream);
return 1;
}
logger_log("INFO", "bridge",
"listening bind_ip=%s listen_port=%u upstream=%s:%u proto=%s client_id=%s",
bind_ip ? bind_ip : "0.0.0.0",
(unsigned)listen_port,
upstream_ip,
(unsigned)upstream_port,
peer_transport_proto_name(proto),
rt.client_id);
while (!g_stop) {
PeerTransportEvent event;
int rc;
rc = peer_transport_next_event(rt.upstream,
&event,
upstream_payload,
sizeof(upstream_payload),
50);
if (rc < 0) {
logger_log("ERROR", "bridge", "upstream_recv_failed rc=%d", rc);
break;
}
if (rc > 0 && handle_upstream_event(&rt, &event, upstream_payload) != OMNI_OK) {
break;
}
rc = peer_transport_next_event(rt.downstream,
&event,
downstream_payload,
sizeof(downstream_payload),
50);
if (rc < 0) {
logger_log("ERROR", "bridge", "downstream_recv_failed rc=%d", rc);
break;
}
if (rc > 0) {
(void)handle_downstream_event(&rt, &event, downstream_payload);
}
}
peer_transport_close(rt.downstream);
peer_transport_close(rt.upstream);
logger_print_performance_log("final");
return 0;
}

601
src/apps/hub_main.c Normal file
View File

@@ -0,0 +1,601 @@
/*
* hub_main.c
* 多客户端 hub:维护 client_id -> session 的映射,并负责 register / bind / tunnel 路由
*
* 支持统一的 tcp / udp / kcp 传输:
* - 多个 peer 主动接入 hub
* - peer 先 REGISTER 自己的逻辑 ID
* - peer 可 BIND 默认目标
* - peer 发送 TUNNEL 后,hub 根据 dst_id 转发给目标
*/
#include "common.h"
#include "logger.h"
#include "peer_transport.h"
#include <ctype.h>
#include <signal.h>
#include <stdint.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#define HUB_MAX_PAYLOAD (PEER_TUNNEL_META_SIZE + 65536u)
typedef struct HubClient {
PeerTransportSession *session;
char client_id[OMNI_PEER_ID_SIZE];
char bound_peer[OMNI_PEER_ID_SIZE];
struct HubClient *next;
} HubClient;
typedef struct HubState {
PeerTransport *transport;
HubClient *clients;
} HubState;
static volatile sig_atomic_t g_stop = 0;
static void on_signal(int signo)
{
(void)signo;
g_stop = 1;
}
static void install_signal_handlers(void)
{
struct sigaction sa;
memset(&sa, 0, sizeof(sa));
sa.sa_handler = on_signal;
sigemptyset(&sa.sa_mask);
(void)sigaction(SIGINT, &sa, NULL);
(void)sigaction(SIGTERM, &sa, NULL);
(void)signal(SIGPIPE, SIG_IGN);
}
static void usage(const char *prog)
{
fprintf(stderr,
"Usage:\n"
" %s -P <listen_port> [-b <bind_ip>] [-p tcp|udp|kcp]\n",
prog);
}
static int parse_proto(const char *s, OmniProtocol *out_proto)
{
if (!s || !out_proto) {
return 0;
}
if (strcmp(s, "tcp") == 0) {
*out_proto = OMNI_PROTO_TCP;
return 1;
}
if (strcmp(s, "udp") == 0) {
*out_proto = OMNI_PROTO_UDP;
return 1;
}
if (strcmp(s, "kcp") == 0) {
*out_proto = OMNI_PROTO_KCP;
return 1;
}
return 0;
}
static int peer_id_is_valid(const char *id)
{
size_t len = 0;
if (!id || !id[0]) {
return 0;
}
for (len = 0; id[len] != '\0'; ++len) {
unsigned char ch = (unsigned char)id[len];
if (len + 1u >= OMNI_PEER_ID_SIZE) {
return 0;
}
if (!(isalnum(ch) || ch == '_' || ch == '-' || ch == '.')) {
return 0;
}
}
return 1;
}
static HubClient *find_client_by_id(HubState *hub, const char *client_id)
{
HubClient *cur;
if (!hub || !client_id) {
return NULL;
}
for (cur = hub->clients; cur; cur = cur->next) {
if (strcmp(cur->client_id, client_id) == 0) {
return cur;
}
}
return NULL;
}
static HubClient *find_client_by_session(HubState *hub, PeerTransportSession *session)
{
HubClient *cur;
if (!hub || !session) {
return NULL;
}
for (cur = hub->clients; cur; cur = cur->next) {
if (cur->session == session) {
return cur;
}
}
return NULL;
}
static int send_status_session(HubState *hub,
PeerTransportSession *session,
uint32_t code,
const char *self_id,
const char *peer_id,
const char *detail)
{
PeerStatusMeta status_meta;
if (!hub || !session) {
return OMNI_ERR_PARAM;
}
omni_peer_status_meta_encode(&status_meta, code, self_id, peer_id, detail);
return peer_transport_send(hub->transport,
session,
MSG_TYPE_PEER_STATUS,
&status_meta,
PEER_STATUS_META_SIZE);
}
static void notify_peer_offline(HubState *hub, const char *departed_id)
{
HubClient *cur;
if (!hub || !departed_id || !departed_id[0]) {
return;
}
for (cur = hub->clients; cur; cur = cur->next) {
if (strcmp(cur->bound_peer, departed_id) == 0) {
cur->bound_peer[0] = '\0';
(void)send_status_session(hub,
cur->session,
PEER_STATUS_UNBOUND,
cur->client_id,
departed_id,
"peer_offline binding_cleared");
}
}
}
static void remove_client(HubState *hub, HubClient *client)
{
HubClient **cur;
char departed_id[OMNI_PEER_ID_SIZE];
if (!hub || !client) {
return;
}
departed_id[0] = '\0';
omni_copy_fixed_ascii(departed_id, sizeof(departed_id), client->client_id);
cur = &hub->clients;
while (*cur) {
if (*cur == client) {
*cur = client->next;
break;
}
cur = &(*cur)->next;
}
if (departed_id[0]) {
notify_peer_offline(hub, departed_id);
}
peer_transport_close_session(hub->transport, client->session);
free(client);
}
static HubClient *register_client(HubState *hub,
PeerTransportSession *session,
const char *client_id)
{
HubClient *client;
HubClient *existing;
char remote_buf[128];
if (!hub || !session || !client_id) {
return NULL;
}
existing = find_client_by_session(hub, session);
if (existing) {
omni_copy_fixed_ascii(existing->client_id, sizeof(existing->client_id), client_id);
return existing;
}
existing = find_client_by_id(hub, client_id);
if (existing && existing->session != session) {
logger_log("WARN", "hub",
"client_id_replaced client_id=%s old_remote=%s",
client_id,
peer_transport_session_remote(existing->session,
remote_buf,
sizeof(remote_buf)));
remove_client(hub, existing);
}
client = (HubClient *)calloc(1, sizeof(*client));
if (!client) {
return NULL;
}
client->session = session;
omni_copy_fixed_ascii(client->client_id, sizeof(client->client_id), client_id);
client->next = hub->clients;
hub->clients = client;
return client;
}
static int handle_register(HubState *hub,
PeerTransportSession *session,
const uint8_t *payload,
uint32_t payload_len)
{
PeerRegisterMeta register_meta;
HubClient *client;
char detail[128];
char remote_buf[128];
if (payload_len < PEER_REGISTER_META_SIZE) {
return send_status_session(hub,
session,
PEER_STATUS_ERROR,
NULL,
NULL,
"short_register_payload");
}
omni_peer_register_meta_decode((const PeerRegisterMeta *)payload, &register_meta);
if (!peer_id_is_valid(register_meta.client_id)) {
return send_status_session(hub,
session,
PEER_STATUS_ERROR,
NULL,
NULL,
"invalid_client_id");
}
client = register_client(hub, session, register_meta.client_id);
if (!client) {
return send_status_session(hub,
session,
PEER_STATUS_ERROR,
NULL,
register_meta.client_id,
"register_alloc_failed");
}
snprintf(detail, sizeof(detail), "registered remote=%s",
peer_transport_session_remote(session, remote_buf, sizeof(remote_buf)));
logger_log("INFO", "hub",
"client_registered client_id=%s remote=%s",
client->client_id,
remote_buf);
return send_status_session(hub,
session,
PEER_STATUS_REGISTERED,
client->client_id,
NULL,
detail);
}
static int handle_bind(HubState *hub,
HubClient *client,
const uint8_t *payload,
uint32_t payload_len)
{
PeerBindMeta bind_meta;
HubClient *target;
if (!hub || !client) {
return OMNI_ERR_PARAM;
}
if (payload_len < PEER_BIND_META_SIZE) {
return send_status_session(hub,
client->session,
PEER_STATUS_ERROR,
client->client_id,
NULL,
"short_bind_payload");
}
omni_peer_bind_meta_decode((const PeerBindMeta *)payload, &bind_meta);
if (!peer_id_is_valid(bind_meta.peer_id)) {
return send_status_session(hub,
client->session,
PEER_STATUS_ERROR,
client->client_id,
NULL,
"invalid_peer_id");
}
if (strcmp(bind_meta.peer_id, client->client_id) == 0) {
return send_status_session(hub,
client->session,
PEER_STATUS_ERROR,
client->client_id,
bind_meta.peer_id,
"cannot_bind_self");
}
target = find_client_by_id(hub, bind_meta.peer_id);
if (!target) {
return send_status_session(hub,
client->session,
PEER_STATUS_ERROR,
client->client_id,
bind_meta.peer_id,
"peer_not_online");
}
omni_copy_fixed_ascii(client->bound_peer, sizeof(client->bound_peer), bind_meta.peer_id);
logger_log("INFO", "hub",
"peer_bound client_id=%s peer_id=%s",
client->client_id,
bind_meta.peer_id);
return send_status_session(hub,
client->session,
PEER_STATUS_BOUND,
client->client_id,
bind_meta.peer_id,
"bind_ok");
}
static int handle_tunnel(HubState *hub,
HubClient *client,
const uint8_t *payload,
uint32_t payload_len)
{
PeerTunnelMeta tunnel_meta;
PeerTunnelMeta forward_meta;
HubClient *target;
char effective_dst[OMNI_PEER_ID_SIZE];
uint8_t *forward_payload = NULL;
uint32_t inner_len;
int rc = OMNI_OK;
if (!hub || !client) {
return OMNI_ERR_PARAM;
}
if (payload_len < PEER_TUNNEL_META_SIZE) {
return send_status_session(hub,
client->session,
PEER_STATUS_ERROR,
client->client_id,
NULL,
"short_tunnel_payload");
}
omni_peer_tunnel_meta_decode((const PeerTunnelMeta *)payload, &tunnel_meta);
inner_len = payload_len - PEER_TUNNEL_META_SIZE;
memset(effective_dst, 0, sizeof(effective_dst));
if (tunnel_meta.dst_id[0] != '\0') {
omni_copy_fixed_ascii(effective_dst, sizeof(effective_dst), tunnel_meta.dst_id);
} else {
omni_copy_fixed_ascii(effective_dst, sizeof(effective_dst), client->bound_peer);
}
if (!peer_id_is_valid(effective_dst)) {
return send_status_session(hub,
client->session,
PEER_STATUS_ERROR,
client->client_id,
NULL,
"missing_or_invalid_destination");
}
target = find_client_by_id(hub, effective_dst);
if (!target) {
return send_status_session(hub,
client->session,
PEER_STATUS_ERROR,
client->client_id,
effective_dst,
"destination_not_online");
}
forward_payload = (uint8_t *)malloc(payload_len);
if (!forward_payload) {
return send_status_session(hub,
client->session,
PEER_STATUS_ERROR,
client->client_id,
effective_dst,
"malloc_forward_payload_failed");
}
omni_peer_tunnel_meta_encode(&forward_meta,
client->client_id,
effective_dst,
tunnel_meta.inner_type);
memcpy(forward_payload, &forward_meta, PEER_TUNNEL_META_SIZE);
if (inner_len > 0) {
memcpy(forward_payload + PEER_TUNNEL_META_SIZE,
payload + PEER_TUNNEL_META_SIZE,
inner_len);
}
rc = peer_transport_send(hub->transport,
target->session,
MSG_TYPE_PEER_TUNNEL,
forward_payload,
payload_len);
free(forward_payload);
if (rc != OMNI_OK) {
logger_log("ERROR", "hub",
"forward_failed src_id=%s dst_id=%s inner_type=%u",
client->client_id,
effective_dst,
(unsigned)tunnel_meta.inner_type);
return send_status_session(hub,
client->session,
PEER_STATUS_ERROR,
client->client_id,
effective_dst,
"forward_failed");
}
logger_log("INFO", "hub",
"forward_ok src_id=%s dst_id=%s inner_type=%u payload_bytes=%u",
client->client_id,
effective_dst,
(unsigned)tunnel_meta.inner_type,
(unsigned)inner_len);
return OMNI_OK;
}
int main(int argc, char **argv)
{
const char *bind_ip = NULL;
const char *proto_str = "tcp";
OmniProtocol proto = OMNI_PROTO_TCP;
int listen_port = 0;
int opt;
HubState hub;
uint8_t payload[HUB_MAX_PAYLOAD];
while ((opt = getopt(argc, argv, "b:p:P:")) != -1) {
switch (opt) {
case 'b':
bind_ip = optarg;
break;
case 'p':
proto_str = optarg;
break;
case 'P':
listen_port = atoi(optarg);
break;
default:
usage(argv[0]);
return 1;
}
}
if (listen_port <= 0 || !parse_proto(proto_str, &proto)) {
usage(argv[0]);
return 1;
}
logger_init();
install_signal_handlers();
logger_set_context("hub", proto_str, "hub", "server", NULL);
memset(&hub, 0, sizeof(hub));
hub.transport = peer_transport_open(OMNI_ROLE_SERVER,
proto,
bind_ip,
(uint16_t)listen_port,
NULL,
0);
if (!hub.transport) {
perror("hub transport");
return 1;
}
logger_log("INFO", "hub",
"listening bind_ip=%s port=%u proto=%s",
bind_ip ? bind_ip : "0.0.0.0",
(unsigned)listen_port,
peer_transport_proto_name(proto));
while (!g_stop) {
PeerTransportEvent event;
HubClient *client;
int rc;
rc = peer_transport_next_event(hub.transport,
&event,
payload,
sizeof(payload),
100);
if (rc < 0) {
logger_log("ERROR", "hub", "transport_recv_failed rc=%d", rc);
continue;
}
if (rc == 0) {
continue;
}
if (event.kind == PEER_TRANSPORT_EVENT_CLOSED) {
client = find_client_by_session(&hub, event.session);
if (client) {
logger_log("INFO", "hub",
"client_closed client_id=%s",
client->client_id[0] ? client->client_id : "unregistered");
remove_client(&hub, client);
} else {
peer_transport_close_session(hub.transport, event.session);
}
continue;
}
switch (event.header.type) {
case MSG_TYPE_PEER_REGISTER:
(void)handle_register(&hub, event.session, payload, event.header.len);
break;
case MSG_TYPE_PEER_BIND:
client = find_client_by_session(&hub, event.session);
if (!client) {
(void)send_status_session(&hub,
event.session,
PEER_STATUS_ERROR,
NULL,
NULL,
"register_first");
break;
}
(void)handle_bind(&hub, client, payload, event.header.len);
break;
case MSG_TYPE_PEER_TUNNEL:
client = find_client_by_session(&hub, event.session);
if (!client) {
(void)send_status_session(&hub,
event.session,
PEER_STATUS_ERROR,
NULL,
NULL,
"register_first");
break;
}
(void)handle_tunnel(&hub, client, payload, event.header.len);
break;
default:
client = find_client_by_session(&hub, event.session);
(void)send_status_session(&hub,
event.session,
PEER_STATUS_ERROR,
client ? client->client_id : NULL,
NULL,
"unsupported_message_type");
break;
}
}
while (hub.clients) {
HubClient *next = hub.clients->next;
peer_transport_close_session(hub.transport, hub.clients->session);
free(hub.clients);
hub.clients = next;
}
peer_transport_close(hub.transport);
logger_print_performance_log("final");
return 0;
}

1966
src/apps/peer_main.c Normal file

File diff suppressed because it is too large Load Diff

File diff suppressed because it is too large Load Diff

View File

@@ -36,6 +36,20 @@ static const struct ProtoVTable *select_vtable(OmniProtocol proto)
} }
} }
static const char *proto_name(OmniProtocol proto)
{
switch (proto) {
case OMNI_PROTO_TCP:
return "tcp";
case OMNI_PROTO_UDP:
return "udp";
case OMNI_PROTO_KCP:
return "kcp";
default:
return "unknown";
}
}
OmniContext *omni_init(OmniRole role, OmniContext *omni_init(OmniRole role,
OmniProtocol proto, OmniProtocol proto,
const char *bind_ip, const char *bind_ip,
@@ -44,6 +58,11 @@ OmniContext *omni_init(OmniRole role,
uint16_t peer_port) uint16_t peer_port)
{ {
logger_init(); logger_init();
logger_set_context("network",
proto_name(proto),
NULL,
role == OMNI_ROLE_SERVER ? "server" : "client",
NULL);
const struct ProtoVTable *vt = select_vtable(proto); const struct ProtoVTable *vt = select_vtable(proto);
if (!vt || !vt->init) { if (!vt || !vt->init) {
@@ -87,12 +106,17 @@ ssize_t omni_send(OmniContext *ctx, const void *buf, size_t len)
return OMNI_ERR_PARAM; return OMNI_ERR_PARAM;
} }
uint64_t t0 = omni_now_ms();
ssize_t n = ctx->vt->send((OmniContext *)ctx->impl, buf, len); ssize_t n = ctx->vt->send((OmniContext *)ctx->impl, buf, len);
uint64_t t1 = omni_now_ms();
logger_on_send_call_latency(t1 - t0);
if (n > 0) { if (n > 0) {
logger_on_send((size_t)n); logger_on_send((size_t)n);
logger_on_send_transmission_bytes((size_t)n);
logger_maybe_print_performance_log("on_send");
} }
logger_log("DEBUG", "network", "omni_send proto=%d bytes=%zd", logger_log("DEBUG", "network", "omni_send proto=%d bytes=%zd call_ms=%llu",
(int)ctx->proto, n); (int)ctx->proto, n, (unsigned long long)(t1 - t0));
return n; return n;
} }
@@ -102,12 +126,17 @@ ssize_t omni_recv(OmniContext *ctx, void *buf, size_t len)
return OMNI_ERR_PARAM; return OMNI_ERR_PARAM;
} }
uint64_t t0 = omni_now_ms();
ssize_t n = ctx->vt->recv((OmniContext *)ctx->impl, buf, len); ssize_t n = ctx->vt->recv((OmniContext *)ctx->impl, buf, len);
uint64_t t1 = omni_now_ms();
logger_on_recv_call_latency(t1 - t0);
if (n > 0) { if (n > 0) {
logger_on_recv((size_t)n); logger_on_recv((size_t)n);
logger_on_recv_transmission_bytes((size_t)n);
logger_maybe_print_performance_log("on_recv");
} }
logger_log("DEBUG", "network", "omni_recv proto=%d bytes=%zd", logger_log("DEBUG", "network", "omni_recv proto=%d bytes=%zd call_ms=%llu",
(int)ctx->proto, n); (int)ctx->proto, n, (unsigned long long)(t1 - t0));
return n; return n;
} }
@@ -120,4 +149,3 @@ void omni_close(OmniContext *ctx)
logger_print_performance_log("final"); logger_print_performance_log("final");
free(ctx); free(ctx);
} }

1432
src/core/peer_transport.c Normal file

File diff suppressed because it is too large Load Diff

View File

@@ -22,8 +22,75 @@ struct KcpContext {
struct sockaddr_in peer_addr; struct sockaddr_in peer_addr;
socklen_t peer_len; socklen_t peer_len;
ikcpcb *kcp; ikcpcb *kcp;
uint32_t *seg_xmit_seen; /* 按 sn 记录上次采样到的 xmit 值 */
size_t seg_xmit_cap;
}; };
static int kcp_ensure_seg_track_capacity(struct KcpContext *ctx, uint32_t sn)
{
size_t need;
size_t new_cap;
uint32_t *new_seen;
if (!ctx) return OMNI_ERR_PARAM;
need = (size_t)sn + 1u;
if (need <= ctx->seg_xmit_cap) {
return OMNI_OK;
}
new_cap = (ctx->seg_xmit_cap == 0) ? 64u : ctx->seg_xmit_cap;
while (new_cap < need) {
new_cap *= 2u;
}
new_seen = (uint32_t *)realloc(ctx->seg_xmit_seen, new_cap * sizeof(uint32_t));
if (!new_seen) {
return OMNI_ERR_GENERIC;
}
memset(new_seen + ctx->seg_xmit_cap, 0,
(new_cap - ctx->seg_xmit_cap) * sizeof(uint32_t));
ctx->seg_xmit_seen = new_seen;
ctx->seg_xmit_cap = new_cap;
return OMNI_OK;
}
static void kcp_sample_transport_stats(struct KcpContext *ctx)
{
struct IQUEUEHEAD *p;
if (!ctx || !ctx->kcp) {
return;
}
for (p = ctx->kcp->snd_buf.next; p != &ctx->kcp->snd_buf; p = p->next) {
struct IKCPSEG *segment = iqueue_entry(p, struct IKCPSEG, node);
uint32_t prev_xmit;
if (!segment || segment->len == 0) {
continue;
}
if (kcp_ensure_seg_track_capacity(ctx, segment->sn) != OMNI_OK) {
logger_log("WARN", "kcp", "seg_track_alloc_failed sn=%u",
(unsigned)segment->sn);
return;
}
prev_xmit = ctx->seg_xmit_seen[segment->sn];
if (prev_xmit == 0 && segment->xmit > 0) {
logger_on_kcp_tx(1u, (uint64_t)segment->len);
prev_xmit = 1u;
}
if (segment->xmit > prev_xmit) {
uint32_t delta = segment->xmit - prev_xmit;
logger_on_kcp_retrans((uint64_t)delta,
(uint64_t)delta * (uint64_t)segment->len);
}
ctx->seg_xmit_seen[segment->sn] = segment->xmit;
}
}
static int kcp_output(const char *buf, int len, ikcpcb *kcp, void *user) static int kcp_output(const char *buf, int len, ikcpcb *kcp, void *user)
{ {
(void)kcp; (void)kcp;
@@ -43,8 +110,6 @@ static OmniContext *kcp_init(OmniRole role,
const char *peer_ip, const char *peer_ip,
uint16_t peer_port) uint16_t peer_port)
{ {
(void)role;
struct KcpContext *ctx = (struct KcpContext *)calloc(1, sizeof(*ctx)); struct KcpContext *ctx = (struct KcpContext *)calloc(1, sizeof(*ctx));
if (!ctx) return NULL; if (!ctx) return NULL;
@@ -77,8 +142,9 @@ static OmniContext *kcp_init(OmniRole role,
ctx->fd = fd; ctx->fd = fd;
/* conv 可简单使用端口号 */ /* conv 必须两端一致:server 用 bind_port,client 用 peer_port */
IUINT32 conv = (IUINT32)peer_port; IUINT32 conv = (role == OMNI_ROLE_SERVER) ? (IUINT32)bind_port
: (IUINT32)peer_port;
ikcpcb *kcp = ikcp_create(conv, ctx); ikcpcb *kcp = ikcp_create(conv, ctx);
if (!kcp) { if (!kcp) {
logger_log("ERROR", "kcp", "ikcp_create_failed"); logger_log("ERROR", "kcp", "ikcp_create_failed");
@@ -94,16 +160,18 @@ static OmniContext *kcp_init(OmniRole role,
ikcp_wndsize(kcp, 128, 128); ikcp_wndsize(kcp, 128, 128);
logger_log("INFO", "kcp", logger_log("INFO", "kcp",
"init bind_port=%u peer_ip=%s peer_port=%u", "init bind_port=%u peer_ip=%s peer_port=%u conv=%u",
(unsigned)bind_port, (unsigned)bind_port,
peer_ip ? peer_ip : "NULL", peer_ip ? peer_ip : "NULL",
(unsigned)peer_port); (unsigned)peer_port,
(unsigned)conv);
return (OmniContext *)ctx; return (OmniContext *)ctx;
} }
static void kcp_update_loop(struct KcpContext *ctx) static void kcp_update_loop(struct KcpContext *ctx)
{ {
uint64_t t0 = omni_now_ms();
IUINT32 current = (IUINT32)omni_now_ms(); IUINT32 current = (IUINT32)omni_now_ms();
ikcp_update(ctx->kcp, current); ikcp_update(ctx->kcp, current);
@@ -117,6 +185,56 @@ static void kcp_update_loop(struct KcpContext *ctx)
ctx->peer_len = fromlen; ctx->peer_len = fromlen;
ikcp_input(ctx->kcp, buf, (long)n); ikcp_input(ctx->kcp, buf, (long)n);
} }
/* KCP 内部状态监控:按分片 xmit 变化统计首次发送和真实重传。 */
kcp_sample_transport_stats(ctx);
/* rx_srtt 是 KCP 平滑后的 RTT(单位 ms),大于 0 才说明已经有有效采样值。 */
if (ctx->kcp->rx_srtt > 0) {
logger_on_rtt((uint64_t)ctx->kcp->rx_srtt);
}
/* cwnd 是当前拥塞窗口大小,用来观察 KCP 的拥塞控制是否在收缩或增长。 */
logger_on_cwnd((double)ctx->kcp->cwnd);
/* 统计发送/接收窗口的占用百分比,方便观察当前缓冲区压力。 */
{
double send_pct = 0.0;
double recv_pct = 0.0;
/* ikcp_waitsnd() 表示还在发送路径中的分片数量,用发送窗口大小换算成占用率。 */
if (ctx->kcp->snd_wnd > 0) {
send_pct = ((double)ikcp_waitsnd(ctx->kcp) * 100.0) / (double)ctx->kcp->snd_wnd;
logger_on_send_queue_bytes((size_t)ikcp_waitsnd(ctx->kcp) * (size_t)ctx->kcp->mss);
}
/* nrcv_que 是已经收到、等待应用层取走的数据包数量,用接收窗口换算成占用率。 */
if (ctx->kcp->rcv_wnd > 0) {
recv_pct = ((double)ctx->kcp->nrcv_que * 100.0) / (double)ctx->kcp->rcv_wnd;
logger_on_recv_queue_bytes((size_t)ctx->kcp->nrcv_que * (size_t)ctx->kcp->mss);
}
/* 将收发两侧的缓冲区占用情况统一上报给监控模块。 */
logger_on_buffer_status(send_pct, recv_pct);
}
uint64_t t1 = omni_now_ms();
logger_log("DEBUG", "kcp",
"update ms=%llu cwnd=%u ssthresh=%u rmt_wnd=%u snd_wnd=%u rcv_wnd=%u "
"rx_srtt=%u rx_rto=%u nsnd_buf=%u nsnd_que=%u nrcv_buf=%u nrcv_que=%u xmit=%u state=%u",
(unsigned long long)(t1 - t0),
(unsigned)ctx->kcp->cwnd,
(unsigned)ctx->kcp->ssthresh,
(unsigned)ctx->kcp->rmt_wnd,
(unsigned)ctx->kcp->snd_wnd,
(unsigned)ctx->kcp->rcv_wnd,
(unsigned)ctx->kcp->rx_srtt,
(unsigned)ctx->kcp->rx_rto,
(unsigned)ctx->kcp->nsnd_buf,
(unsigned)ctx->kcp->nsnd_que,
(unsigned)ctx->kcp->nrcv_buf,
(unsigned)ctx->kcp->nrcv_que,
(unsigned)ctx->kcp->xmit,
(unsigned)ctx->kcp->state);
} }
static ssize_t kcp_send(OmniContext *c, const void *buf, size_t len) static ssize_t kcp_send(OmniContext *c, const void *buf, size_t len)
@@ -124,6 +242,7 @@ static ssize_t kcp_send(OmniContext *c, const void *buf, size_t len)
struct KcpContext *ctx = (struct KcpContext *)c; struct KcpContext *ctx = (struct KcpContext *)c;
if (!ctx || !ctx->kcp) return OMNI_ERR_PARAM; if (!ctx || !ctx->kcp) return OMNI_ERR_PARAM;
uint64_t t0 = omni_now_ms();
int rc = ikcp_send(ctx->kcp, (const char *)buf, (int)len); int rc = ikcp_send(ctx->kcp, (const char *)buf, (int)len);
if (rc < 0) { if (rc < 0) {
logger_log("ERROR", "kcp", "ikcp_send_failed rc=%d", rc); logger_log("ERROR", "kcp", "ikcp_send_failed rc=%d", rc);
@@ -132,6 +251,10 @@ static ssize_t kcp_send(OmniContext *c, const void *buf, size_t len)
/* 驱动一次 flush */ /* 驱动一次 flush */
kcp_update_loop(ctx); kcp_update_loop(ctx);
uint64_t t1 = omni_now_ms();
logger_on_proto_send_latency(t1 - t0);
logger_log("DEBUG", "kcp", "send payload_bytes=%zu proto_ms=%llu waitsnd=%d",
len, (unsigned long long)(t1 - t0), ikcp_waitsnd(ctx->kcp));
return (ssize_t)len; return (ssize_t)len;
} }
@@ -140,12 +263,17 @@ static ssize_t kcp_recv(OmniContext *c, void *buf, size_t len)
struct KcpContext *ctx = (struct KcpContext *)c; struct KcpContext *ctx = (struct KcpContext *)c;
if (!ctx || !ctx->kcp) return OMNI_ERR_PARAM; if (!ctx || !ctx->kcp) return OMNI_ERR_PARAM;
uint64_t t0 = omni_now_ms();
kcp_update_loop(ctx); kcp_update_loop(ctx);
int n = ikcp_recv(ctx->kcp, (char *)buf, (int)len); int n = ikcp_recv(ctx->kcp, (char *)buf, (int)len);
if (n < 0) { if (n < 0) {
return 0; /* 暂无数据 */ return 0; /* 暂无数据 */
} }
uint64_t t1 = omni_now_ms();
logger_on_proto_recv_latency(t1 - t0);
logger_log("DEBUG", "kcp", "recv payload_bytes=%d proto_ms=%llu",
n, (unsigned long long)(t1 - t0));
return (ssize_t)n; return (ssize_t)n;
} }
@@ -159,6 +287,7 @@ static void kcp_close(OmniContext *c)
if (ctx->fd >= 0) { if (ctx->fd >= 0) {
close(ctx->fd); close(ctx->fd);
} }
free(ctx->seg_xmit_seen);
free(ctx); free(ctx);
} }

View File

@@ -1,6 +1,12 @@
/* /*
* tcp_impl.c * tcp_impl.c
* TCP 协议实现,带 16 字节包头解决粘包 * TCP 协议实现,带 16 字节包头解决粘包
*
* 设计说明:
* 1) TCP 是字节流,天然没有消息边界,因此这里通过“固定 16 字节头 + payload 长度”
* 显式划分消息边界,避免粘包/拆包带来的上层读取混乱。
* 2) 本层的头只用于“流边界管理”,上层业务仍可在 payload 中定义自己的消息头。
* 3) send/recv 均采用阻塞全量读写语义:要么完整收发一帧,要么返回错误/关闭状态。
*/ */
#include "common.h" #include "common.h"
@@ -10,25 +16,236 @@
#include <arpa/inet.h> #include <arpa/inet.h>
#include <errno.h> #include <errno.h>
#include <netinet/tcp.h> #include <netinet/tcp.h>
#include <netinet/in.h>
#include <stddef.h>
#include <stdio.h> #include <stdio.h>
#include <stdlib.h> #include <stdlib.h>
#include <string.h> #include <string.h>
#include <sys/ioctl.h>
#include <sys/socket.h> #include <sys/socket.h>
#include <sys/types.h> #include <sys/types.h>
#include <unistd.h> #include <unistd.h>
/*
* Linux 下 glibc 的 <netinet/tcp.h> 只暴露了 tcp_info 的旧前缀字段,
* 这里复制到 tcpi_bytes_retrans 为止的内核布局,用来安全读取扩展快照。
*/
#ifdef __linux__
struct OmniLinuxTcpInfo {
uint8_t tcpi_state;
uint8_t tcpi_ca_state;
uint8_t tcpi_retransmits;
uint8_t tcpi_probes;
uint8_t tcpi_backoff;
uint8_t tcpi_options;
uint8_t tcpi_snd_wscale : 4, tcpi_rcv_wscale : 4;
uint8_t tcpi_delivery_rate_app_limited : 1, tcpi_fastopen_client_fail : 2;
uint32_t tcpi_rto;
uint32_t tcpi_ato;
uint32_t tcpi_snd_mss;
uint32_t tcpi_rcv_mss;
uint32_t tcpi_unacked;
uint32_t tcpi_sacked;
uint32_t tcpi_lost;
uint32_t tcpi_retrans;
uint32_t tcpi_fackets;
uint32_t tcpi_last_data_sent;
uint32_t tcpi_last_ack_sent;
uint32_t tcpi_last_data_recv;
uint32_t tcpi_last_ack_recv;
uint32_t tcpi_pmtu;
uint32_t tcpi_rcv_ssthresh;
uint32_t tcpi_rtt;
uint32_t tcpi_rttvar;
uint32_t tcpi_snd_ssthresh;
uint32_t tcpi_snd_cwnd;
uint32_t tcpi_advmss;
uint32_t tcpi_reordering;
uint32_t tcpi_rcv_rtt;
uint32_t tcpi_rcv_space;
uint32_t tcpi_total_retrans;
uint64_t tcpi_pacing_rate;
uint64_t tcpi_max_pacing_rate;
uint64_t tcpi_bytes_acked;
uint64_t tcpi_bytes_received;
uint32_t tcpi_segs_out;
uint32_t tcpi_segs_in;
uint32_t tcpi_notsent_bytes;
uint32_t tcpi_min_rtt;
uint32_t tcpi_data_segs_in;
uint32_t tcpi_data_segs_out;
uint64_t tcpi_delivery_rate;
uint64_t tcpi_busy_time;
uint64_t tcpi_rwnd_limited;
uint64_t tcpi_sndbuf_limited;
uint32_t tcpi_delivered;
uint32_t tcpi_delivered_ce;
uint64_t tcpi_bytes_sent;
uint64_t tcpi_bytes_retrans;
};
static int tcp_info_has_field(socklen_t len, size_t field_end)
{
return (size_t)len >= field_end;
}
static uint64_t tcp_info_rtt_us_to_ms(uint32_t rtt_us)
{
if (rtt_us == 0u) {
return 0u;
}
return ((uint64_t)rtt_us + 999u) / 1000u;
}
#endif
struct TcpContext { struct TcpContext {
/* 已建立连接的 socket fd(服务端 accept 后或客户端 connect 后)。 */
int fd; int fd;
}; };
#ifdef __linux__
static void tcp_sample_socket_buffers(int fd)
{
/* socket 层配置的发送/接收缓冲区总大小。 */
int sndbuf = 0;
int rcvbuf = 0;
socklen_t optlen = sizeof(int);
/* 当前内核发送队列/接收队列里还压着的字节数。 */
int outq = 0;
int inq = 0;
/* 最终换算成百分比后上报,表示缓冲区占用压力。 */
double send_pct = 0.0;
double recv_pct = 0.0;
/*
* SO_SNDBUF 取发送缓冲区容量;
* TIOCOUTQ 取当前还没真正发出去的字节数。
* 两者相除后得到发送侧缓冲占用率。
*/
if (getsockopt(fd, SOL_SOCKET, SO_SNDBUF, &sndbuf, &optlen) == 0 && sndbuf > 0) {
#ifdef TIOCOUTQ
if (ioctl(fd, TIOCOUTQ, &outq) == 0 && outq >= 0) {
send_pct = ((double)outq * 100.0) / (double)sndbuf;
logger_on_send_queue_bytes((size_t)outq);
}
#endif
}
/*
* SO_RCVBUF 取接收缓冲区容量;
* FIONREAD 取当前已经到达、但应用层还没 read 的字节数。
* 两者相除后得到接收侧缓冲占用率。
*/
optlen = sizeof(int);
if (getsockopt(fd, SOL_SOCKET, SO_RCVBUF, &rcvbuf, &optlen) == 0 && rcvbuf > 0) {
if (ioctl(fd, FIONREAD, &inq) == 0 && inq >= 0) {
recv_pct = ((double)inq * 100.0) / (double)rcvbuf;
logger_on_recv_queue_bytes((size_t)inq);
}
}
/* 将 TCP 收发缓冲区的占用情况统一交给 logger 记录。 */
logger_on_buffer_status(send_pct, recv_pct);
}
static void tcp_log_info(int fd, const char *tag)
{
struct OmniLinuxTcpInfo ti;
uint64_t total_retrans = 0;
uint64_t data_segs_out = 0;
uint64_t bytes_sent = 0;
uint64_t bytes_retrans = 0;
socklen_t len = sizeof(ti);
memset(&ti, 0, sizeof(ti));
if (getsockopt(fd, IPPROTO_TCP, TCP_INFO, &ti, &len) != 0) {
return;
}
/* 注意:tcpi_rtt / tcpi_rttvar 单位通常为微秒(Linux),这里向上取整到 ms。 */
unsigned long long rtt_ms = (unsigned long long)tcp_info_rtt_us_to_ms(ti.tcpi_rtt);
unsigned long long rttvar_ms = (unsigned long long)tcp_info_rtt_us_to_ms(ti.tcpi_rttvar);
if (tcp_info_has_field(len, offsetof(struct OmniLinuxTcpInfo, tcpi_total_retrans) +
sizeof(ti.tcpi_total_retrans))) {
total_retrans = (uint64_t)ti.tcpi_total_retrans;
} else {
total_retrans = (uint64_t)ti.tcpi_retrans;
}
if (tcp_info_has_field(len, offsetof(struct OmniLinuxTcpInfo, tcpi_data_segs_out) +
sizeof(ti.tcpi_data_segs_out))) {
data_segs_out = (uint64_t)ti.tcpi_data_segs_out;
}
if (tcp_info_has_field(len, offsetof(struct OmniLinuxTcpInfo, tcpi_bytes_sent) +
sizeof(ti.tcpi_bytes_sent))) {
bytes_sent = (uint64_t)ti.tcpi_bytes_sent;
}
if (tcp_info_has_field(len, offsetof(struct OmniLinuxTcpInfo, tcpi_bytes_retrans) +
sizeof(ti.tcpi_bytes_retrans))) {
bytes_retrans = (uint64_t)ti.tcpi_bytes_retrans;
}
logger_log("INFO", "tcpinfo",
"tag=%s state=%u retransmits=%u probes=%u backoff=%u "
"rto=%u ato=%u rtt_ms=%llu rttvar_ms=%llu "
"snd_cwnd=%u snd_ssthresh=%u snd_mss=%u rcv_mss=%u "
"lost=%u retrans=%u total_retrans=%llu data_segs_out=%llu "
"bytes_sent=%llu bytes_retrans=%llu fackets=%u "
"last_data_sent_ms=%u last_data_recv_ms=%u",
tag ? tag : "sample",
(unsigned)ti.tcpi_state,
(unsigned)ti.tcpi_retransmits,
(unsigned)ti.tcpi_probes,
(unsigned)ti.tcpi_backoff,
(unsigned)ti.tcpi_rto,
(unsigned)ti.tcpi_ato,
rtt_ms,
rttvar_ms,
(unsigned)ti.tcpi_snd_cwnd,
(unsigned)ti.tcpi_snd_ssthresh,
(unsigned)ti.tcpi_snd_mss,
(unsigned)ti.tcpi_rcv_mss,
(unsigned)ti.tcpi_lost,
(unsigned)ti.tcpi_retrans,
(unsigned long long)total_retrans,
(unsigned long long)data_segs_out,
(unsigned long long)bytes_sent,
(unsigned long long)bytes_retrans,
(unsigned)ti.tcpi_fackets,
(unsigned)ti.tcpi_last_data_sent,
(unsigned)ti.tcpi_last_data_recv);
logger_on_rtt(rtt_ms);
logger_on_tcp_transport(total_retrans, data_segs_out, bytes_sent, bytes_retrans);
logger_on_cwnd((double)ti.tcpi_snd_cwnd);
tcp_sample_socket_buffers(fd);
}
#endif
static int tcp_set_nodelay(int fd) static int tcp_set_nodelay(int fd)
{ {
/* 关闭 Nagle,降低小包时延(更利于交互指令场景)。 */
int flag = 1; int flag = 1;
return setsockopt(fd, IPPROTO_TCP, TCP_NODELAY, &flag, sizeof(flag)); return setsockopt(fd, IPPROTO_TCP, TCP_NODELAY, &flag, sizeof(flag));
} }
static int tcp_set_reuseaddr(int fd) static int tcp_set_reuseaddr(int fd)
{ {
/* 允许端口快速复用,减少开发/测试时 TIME_WAIT 影响。 */
int flag = 1; int flag = 1;
return setsockopt(fd, SOL_SOCKET, SO_REUSEADDR, &flag, sizeof(flag)); return setsockopt(fd, SOL_SOCKET, SO_REUSEADDR, &flag, sizeof(flag));
} }
@@ -37,12 +254,14 @@ static int tcp_bind_and_listen(struct TcpContext *ctx,
const char *bind_ip, const char *bind_ip,
uint16_t bind_port) uint16_t bind_port)
{ {
/* 创建监听 socket。 */
int fd = socket(AF_INET, SOCK_STREAM, 0); int fd = socket(AF_INET, SOCK_STREAM, 0);
if (fd < 0) { if (fd < 0) {
logger_log("ERROR", "tcp", "socket_failed errno=%d", errno); logger_log("ERROR", "tcp", "socket_failed errno=%d", errno);
return -1; return -1;
} }
/* 监听 socket 打开地址复用。 */
tcp_set_reuseaddr(fd); tcp_set_reuseaddr(fd);
struct sockaddr_in addr; struct sockaddr_in addr;
@@ -57,6 +276,10 @@ static int tcp_bind_and_listen(struct TcpContext *ctx,
return -1; return -1;
} }
/*
* 这里 backlog 取 1,符合当前“单连接演示/测试”场景。
* 若后续要支持多客户端,可提升 backlog 并改为事件循环/线程池模型。
*/
if (listen(fd, 1) < 0) { if (listen(fd, 1) < 0) {
logger_log("ERROR", "tcp", "listen_failed errno=%d", errno); logger_log("ERROR", "tcp", "listen_failed errno=%d", errno);
close(fd); close(fd);
@@ -65,7 +288,7 @@ static int tcp_bind_and_listen(struct TcpContext *ctx,
logger_log("INFO", "tcp", "listening port=%u", (unsigned)bind_port); logger_log("INFO", "tcp", "listening port=%u", (unsigned)bind_port);
/* 简化:阻塞接受一个客户端,之后用于长连接 */ /* 简化:阻塞接受一个客户端,连接建立后作为长连接使用。 */
int cfd = accept(fd, NULL, NULL); int cfd = accept(fd, NULL, NULL);
if (cfd < 0) { if (cfd < 0) {
logger_log("ERROR", "tcp", "accept_failed errno=%d", errno); logger_log("ERROR", "tcp", "accept_failed errno=%d", errno);
@@ -73,6 +296,7 @@ static int tcp_bind_and_listen(struct TcpContext *ctx,
return -1; return -1;
} }
/* 监听 fd 仅用于 accept,一旦接入成功即可关闭监听 fd。 */
close(fd); close(fd);
tcp_set_nodelay(cfd); tcp_set_nodelay(cfd);
@@ -84,6 +308,7 @@ static int tcp_connect_peer(struct TcpContext *ctx,
const char *peer_ip, const char *peer_ip,
uint16_t peer_port) uint16_t peer_port)
{ {
/* 创建主动连接 socket。 */
int fd = socket(AF_INET, SOCK_STREAM, 0); int fd = socket(AF_INET, SOCK_STREAM, 0);
if (fd < 0) { if (fd < 0) {
logger_log("ERROR", "tcp", "socket_failed errno=%d", errno); logger_log("ERROR", "tcp", "socket_failed errno=%d", errno);
@@ -96,6 +321,7 @@ static int tcp_connect_peer(struct TcpContext *ctx,
addr.sin_port = htons(peer_port); addr.sin_port = htons(peer_port);
addr.sin_addr.s_addr = inet_addr(peer_ip); addr.sin_addr.s_addr = inet_addr(peer_ip);
/* 阻塞 connect 到对端。 */
if (connect(fd, (struct sockaddr *)&addr, sizeof(addr)) < 0) { if (connect(fd, (struct sockaddr *)&addr, sizeof(addr)) < 0) {
logger_log("ERROR", "tcp", "connect_failed errno=%d", errno); logger_log("ERROR", "tcp", "connect_failed errno=%d", errno);
close(fd); close(fd);
@@ -111,6 +337,12 @@ static int tcp_connect_peer(struct TcpContext *ctx,
static ssize_t tcp_read_n(int fd, void *buf, size_t n) static ssize_t tcp_read_n(int fd, void *buf, size_t n)
{ {
/*
* 从 TCP 流中“恰好读取 n 字节”:
* - 正常返回 n
* - 返回 0 表示对端关闭(如果发生在中途,返回已读字节数)
* - 返回 -1 表示系统调用错误
*/
size_t off = 0; size_t off = 0;
char *p = (char *)buf; char *p = (char *)buf;
while (off < n) { while (off < n) {
@@ -129,6 +361,11 @@ static ssize_t tcp_read_n(int fd, void *buf, size_t n)
static ssize_t tcp_write_n(int fd, const void *buf, size_t n) static ssize_t tcp_write_n(int fd, const void *buf, size_t n)
{ {
/*
* 向 TCP 流中“恰好写入 n 字节”:
* - EINTR 自动重试
* - 其余错误返回 -1
*/
size_t off = 0; size_t off = 0;
const char *p = (const char *)buf; const char *p = (const char *)buf;
while (off < n) { while (off < n) {
@@ -148,12 +385,14 @@ static OmniContext *tcp_init(OmniRole role,
const char *peer_ip, const char *peer_ip,
uint16_t peer_port) uint16_t peer_port)
{ {
/* 协议私有上下文(通过 OmniContext* 向上层做不透明传递)。 */
struct TcpContext *ctx = (struct TcpContext *)calloc(1, sizeof(*ctx)); struct TcpContext *ctx = (struct TcpContext *)calloc(1, sizeof(*ctx));
if (!ctx) { if (!ctx) {
return NULL; return NULL;
} }
int rc; int rc;
/* 按角色决定是被动监听还是主动连接。 */
if (role == OMNI_ROLE_SERVER) { if (role == OMNI_ROLE_SERVER) {
rc = tcp_bind_and_listen(ctx, bind_ip, bind_port); rc = tcp_bind_and_listen(ctx, bind_ip, bind_port);
} else { } else {
@@ -178,14 +417,18 @@ static ssize_t tcp_send(OmniContext *c, const void *buf, size_t len)
struct TcpContext *ctx = (struct TcpContext *)c; struct TcpContext *ctx = (struct TcpContext *)c;
if (!ctx || ctx->fd < 0) return OMNI_ERR_PARAM; if (!ctx || ctx->fd < 0) return OMNI_ERR_PARAM;
/*
* 外层 TCP 帧头(16B)仅用于切分消息边界。
* 当前 type 统一标记为 MSG_TYPE_RAW,表示“payload 是上层透传内容”。
*/
uint64_t t0 = omni_now_ms();
MsgHeader hdr; MsgHeader hdr;
hdr.magic = htonl(MSG_MAGIC); omni_msg_header_encode(&hdr, MSG_TYPE_RAW, (uint32_t)len, t0);
hdr.length = htonl((uint32_t)len);
hdr.seq = 0; /* 如有需要,上层可扩展维护序列号 */
uint8_t header_buf[MSG_HEADER_SIZE]; uint8_t header_buf[MSG_HEADER_SIZE];
memcpy(header_buf, &hdr, MSG_HEADER_SIZE); memcpy(header_buf, &hdr, MSG_HEADER_SIZE);
/* 先写固定头,再写 payload,接收侧可据此恢复完整帧。 */
ssize_t n1 = tcp_write_n(ctx->fd, header_buf, MSG_HEADER_SIZE); ssize_t n1 = tcp_write_n(ctx->fd, header_buf, MSG_HEADER_SIZE);
if (n1 != (ssize_t)MSG_HEADER_SIZE) { if (n1 != (ssize_t)MSG_HEADER_SIZE) {
return OMNI_ERR_IO; return OMNI_ERR_IO;
@@ -196,6 +439,14 @@ static ssize_t tcp_send(OmniContext *c, const void *buf, size_t len)
return OMNI_ERR_IO; return OMNI_ERR_IO;
} }
uint64_t t1 = omni_now_ms();
/* 记录协议层发送耗时,便于后续性能分析。 */
logger_on_proto_send_latency(t1 - t0);
logger_log("DEBUG", "tcp", "send payload_bytes=%zu header_bytes=%zu proto_ms=%llu",
len, (size_t)MSG_HEADER_SIZE, (unsigned long long)(t1 - t0));
#ifdef __linux__
tcp_log_info(ctx->fd, "after_send");
#endif
return (ssize_t)len; return (ssize_t)len;
} }
@@ -204,6 +455,13 @@ static ssize_t tcp_recv(OmniContext *c, void *buf, size_t len)
struct TcpContext *ctx = (struct TcpContext *)c; struct TcpContext *ctx = (struct TcpContext *)c;
if (!ctx || ctx->fd < 0) return OMNI_ERR_PARAM; if (!ctx || ctx->fd < 0) return OMNI_ERR_PARAM;
/*
* 收包流程:
* 1) 固定先读 16 字节头
* 2) 解析 payload_len
* 3) 再读 payload_len 字节
*/
uint64_t t0 = omni_now_ms();
uint8_t header_buf[MSG_HEADER_SIZE]; uint8_t header_buf[MSG_HEADER_SIZE];
ssize_t n1 = tcp_read_n(ctx->fd, header_buf, MSG_HEADER_SIZE); ssize_t n1 = tcp_read_n(ctx->fd, header_buf, MSG_HEADER_SIZE);
if (n1 <= 0) { if (n1 <= 0) {
@@ -213,14 +471,17 @@ static ssize_t tcp_recv(OmniContext *c, void *buf, size_t len)
return OMNI_ERR_IO; return OMNI_ERR_IO;
} }
/* 解码网络字节序头字段。 */
MsgHeader hdr; MsgHeader hdr;
MsgHeader host_hdr;
memcpy(&hdr, header_buf, MSG_HEADER_SIZE); memcpy(&hdr, header_buf, MSG_HEADER_SIZE);
if (ntohl(hdr.magic) != MSG_MAGIC) { omni_msg_header_decode(&hdr, &host_hdr);
logger_log("ERROR", "tcp", "invalid_magic");
return OMNI_ERR_IO;
}
uint32_t payload_len = ntohl(hdr.length); uint32_t payload_len = host_hdr.len;
/*
* 调用方缓冲区不足时直接报错。
* 当前实现不做“读取并丢弃剩余字节”,因此调用方应保证 recv 缓冲足够大。
*/
if (payload_len > len) { if (payload_len > len) {
logger_log("ERROR", "tcp", "buffer_too_small payload=%u buf_len=%zu", logger_log("ERROR", "tcp", "buffer_too_small payload=%u buf_len=%zu",
payload_len, len); payload_len, len);
@@ -233,6 +494,18 @@ static ssize_t tcp_recv(OmniContext *c, void *buf, size_t len)
return OMNI_ERR_IO; return OMNI_ERR_IO;
} }
uint64_t t1 = omni_now_ms();
/* 记录协议层接收耗时。 */
logger_on_proto_recv_latency(t1 - t0);
logger_log("DEBUG", "tcp",
"recv payload_bytes=%u header_bytes=%zu msg_type=%u ts_ms=%llu proto_ms=%llu",
payload_len, (size_t)MSG_HEADER_SIZE,
(unsigned)host_hdr.type,
(unsigned long long)host_hdr.timestamp,
(unsigned long long)(t1 - t0));
#ifdef __linux__
tcp_log_info(ctx->fd, "after_recv");
#endif
return (ssize_t)payload_len; return (ssize_t)payload_len;
} }
@@ -240,6 +513,7 @@ static void tcp_close(OmniContext *c)
{ {
struct TcpContext *ctx = (struct TcpContext *)c; struct TcpContext *ctx = (struct TcpContext *)c;
if (!ctx) return; if (!ctx) return;
/* 关闭连接并释放私有上下文。 */
if (ctx->fd >= 0) { if (ctx->fd >= 0) {
close(ctx->fd); close(ctx->fd);
} }

View File

@@ -12,6 +12,7 @@
#include <stdio.h> #include <stdio.h>
#include <stdlib.h> #include <stdlib.h>
#include <string.h> #include <string.h>
#include <sys/ioctl.h>
#include <sys/socket.h> #include <sys/socket.h>
#include <sys/types.h> #include <sys/types.h>
#include <unistd.h> #include <unistd.h>
@@ -22,6 +23,52 @@ struct UdpContext {
socklen_t peer_len; socklen_t peer_len;
}; };
static void udp_sample_socket_buffers(int fd)
{
/* socket 层配置的发送/接收缓冲区总大小。 */
int sndbuf = 0;
int rcvbuf = 0;
socklen_t optlen = sizeof(int);
/* 当前内核发送队列/接收队列里还压着的字节数。 */
int outq = 0;
int inq = 0;
/* 最终换算成百分比后上报,表示缓冲区占用压力。 */
double send_pct = 0.0;
double recv_pct = 0.0;
/*
* SO_SNDBUF 取发送缓冲区容量;
* TIOCOUTQ 取当前还没从内核发送队列发走的字节数。
* 两者相除后得到发送侧缓冲占用率。
*/
if (getsockopt(fd, SOL_SOCKET, SO_SNDBUF, &sndbuf, &optlen) == 0 && sndbuf > 0) {
#ifdef TIOCOUTQ
if (ioctl(fd, TIOCOUTQ, &outq) == 0 && outq >= 0) {
send_pct = ((double)outq * 100.0) / (double)sndbuf;
logger_on_send_queue_bytes((size_t)outq);
}
#endif
}
/*
* SO_RCVBUF 取接收缓冲区容量;
* FIONREAD 取当前已经到达、但应用层还没 recv 的字节数。
* 两者相除后得到接收侧缓冲占用率。
*/
optlen = sizeof(int);
if (getsockopt(fd, SOL_SOCKET, SO_RCVBUF, &rcvbuf, &optlen) == 0 && rcvbuf > 0) {
if (ioctl(fd, FIONREAD, &inq) == 0 && inq >= 0) {
recv_pct = ((double)inq * 100.0) / (double)rcvbuf;
logger_on_recv_queue_bytes((size_t)inq);
}
}
/* 将 UDP 收发缓冲区的占用情况统一交给 logger 记录。 */
logger_on_buffer_status(send_pct, recv_pct);
}
static OmniContext *udp_init(OmniRole role, static OmniContext *udp_init(OmniRole role,
const char *bind_ip, const char *bind_ip,
uint16_t bind_port, uint16_t bind_port,
@@ -82,6 +129,7 @@ static ssize_t udp_send(OmniContext *c, const void *buf, size_t len)
logger_log("ERROR", "udp", "sendto_failed errno=%d", errno); logger_log("ERROR", "udp", "sendto_failed errno=%d", errno);
return OMNI_ERR_IO; return OMNI_ERR_IO;
} }
udp_sample_socket_buffers(ctx->fd);
return n; return n;
} }
@@ -103,6 +151,7 @@ static ssize_t udp_recv(OmniContext *c, void *buf, size_t len)
/* 默认更新 peer 为最近一次通信对端,便于“伪长连接” */ /* 默认更新 peer 为最近一次通信对端,便于“伪长连接” */
ctx->peer_addr = from; ctx->peer_addr = from;
ctx->peer_len = fromlen; ctx->peer_len = fromlen;
udp_sample_socket_buffers(ctx->fd);
return n; return n;
} }

Binary file not shown.

Binary file not shown.