WuKongIM Docs

Rust 快速接入

通过 WuKongEasySDK-Rust 和 Tokio 连接用户,完成在线收发、有限重连和连接清理。

编辑此页报告文档问题

WuKongEasySDK-Rust 适合原生 Rust 程序、桌面应用和已有业务后端的通信客户端。它参考 WuKongEasySDK-JS 2.0.4,使用 WebSocket JSON-RPC CONNECT 完成鉴权,再收发在线消息。

当前 Product Gateway 支持这条在线双向收发路径。SDK 不提供本地消息库、会话/未读、离线恢复或推送;需要这些能力时先阅读 SDK 选择

1. 安装正式版本

官方仓库:WuKongIM/WuKongEasySDK-Rust。已发布 crates.io 0.1.0,需要 Rust 1.86+ 和 Tokio。使用以下精确版本:

[dependencies]
wukong-easy-sdk = "=0.1.0"
tokio = { version = "1", features = ["macros", "rt-multi-thread"] }
serde_json = "1"

包名是 wukong-easy-sdk,Rust 导入名是 wukong_easy_sdk。提交应用的 Cargo.lock,让后续构建复用已解析的依赖。当前实现支持原生 TCP/TLS,WSS 使用 rustls 与 WebPKI 根证书;不支持浏览器/WASM。

2. 获取连接材料

业务后端返回当前用户的 uid、短期 tokenwebsocketUrl。按认证与 Token保护注册和轮换流程;客户端不应直接调用 Product HTTP 管理接口。

Auth::new 默认设备类别是 PC/Desktop 2。协议值为 APP 0、WEB 1、PC 2,Token 必须按同一设备类别保存。若你的原生宿主使用 APP 类别,显式设置 auth.device_flag = DeviceFlag::App,并同步后端注册类别。

Auth::new 生成的 device_id 在同一个客户端重连时保持不变;需要跨进程保留设备身份时,由应用读取并设置自己的稳定设备 ID。每个身份创建一个 Clientclone() 共享同一条连接,适合在 Tokio 任务之间传递。

3. 连接、发送并清理

这个完整程序先订阅事件,再连接并向 Bob 发一条消息,随后清理。把环境变量换成业务后端提供的材料。Alice 和 Bob 必须是不同 UID,且 Bob 已在线。

use serde_json::json;
use wukong_easy_sdk::{Auth, ChannelType, Client, Event, Options};

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let client = Client::new(
        std::env::var("WK_WS_URL")?,
        Auth::new(std::env::var("WK_UID")?, std::env::var("WK_TOKEN")?),
        Options::default(),
    )?;
    let mut events = client.subscribe();
    let listener = tokio::spawn(async move {
        loop {
            match events.recv().await {
                Ok(Event::Message(message)) => {
                    // 在应用 UI 中展示或保存 message.payload,不记录完整正文。
                    let _ = &message.payload;
                }
                Ok(Event::CustomEvent(event)) => {
                    // event.event_type 和 event.data 由业务协议定义。
                    let _ = (&event.event_type, &event.data);
                }
                Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => {
                    // 已丢失事件;通知应用通过业务后端补偿状态。
                    break;
                }
                Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
                _ => {}
            }
        }
    });
    let result = async {
        client.connect().await?;
        let ack = client.send(
            "bob", ChannelType::Person,
            json!({"type": 1, "content": "你好,Rust 🦀"}),
        ).await?;
        // ack 是服务端发送结果,不是 Bob 已读或业务处理完成。
        let _ = ack;
        Ok::<_, wukong_easy_sdk::Error>(())
    }.await;
    client.destroy().await;
    listener.abort();
    let _ = listener.await;
    result?;
    Ok(())
}

成功的 send 返回 SendResult,包括字符串 message_idu64 类型的 message_seqreason_code。业务错误返回 Error::Server { code },保留服务端数值,不把失败响应当成功。不要将 SENDACK 等同于接收、展示或业务完成,详见消息收发

4. 跑通 Alice 与 Bob

要验证持续双向通信,运行仓库里的终端 example,而不是上面发送后立即退出的程序:

git clone https://github.com/WuKongIM/WuKongEasySDK-Rust.git
cd WuKongEasySDK-Rust
git checkout 5b4a59cdbb66a9e0c3878e73ba4656f08ee05c6b
cargo test --locked

由受信开发终端按运行官方示例准备两组 Token,将 Alice 与 Bob 的 device_flag 都设为 2。在两个终端分别运行:

# Alice
WK_WS_URL=ws://127.0.0.1:5200 WK_UID=alice WK_TOKEN=alice-token \
  WK_PEER_UID=bob cargo run --locked --example chat

# Bob
WK_WS_URL=ws://127.0.0.1:5200 WK_UID=bob WK_TOKEN=bob-token \
  WK_PEER_UID=alice cargo run --locked --example chat

两端都显示连接成功后再输入内容。示例会报告服务端接受发送、收到消息和连接状态,避免把消息正文写入日志;应用可以从 Event::Message 中读取 from_uid、Channel 和 Payload。输入 /quit 或 Ctrl-C 清理退出。

无人值守的双向校验可使用 roundtrip

WK_WS_URL=ws://127.0.0.1:5200 WK_UID=alice WK_TOKEN=alice-token \
  WK_PEER_UID=bob WK_PEER_TOKEN=bob-token \
  cargo run --locked --example roundtrip

该 example 创建两个 Rust 客户端,核对 Unicode Payload 与发送结果,跨多个心跳周期检查连接,并手动重连后重复收发。JS 互通步骤见仓库 README

5. 管理生命周期与资源

任务Rust API 与语义
订阅/取消订阅subscribe() 返回 Tokio broadcast 接收器;drop 接收器取消订阅
并发连接connect().await 共享当前尝试;取消一个 Future 不会取消共享连接
首次失败返回错误,由应用决定是否再次连接
网络重连已建立的连接异常断开后,默认最多 5 次指数退避重连,包含抖动
手动断开disconnect().await 取消鉴权、I/O 和重连;之后可再次 connect
账号退出destroy().await 永久关闭所有 clone,并结束订阅任务
最后一个句柄 drop取消后台连接;需要确认清理完成时显式 await

鉴权拒绝、服务端 disconnect 和手动退出不会自动重连。默认心跳间隔 25 秒,等待 pong 10 秒;连接超时 5 秒,发送总超时 15 秒,写入超时 5 秒。

默认在队列中和等待响应的 SEND 总数不超过 256,超限返回 Backpressure;事件保留 256 条,慢监听器收到 RecvError::Lagged。单个完整 JSON-RPC 报文默认限制 1 MiB,包含 Base64 膨胀。通过 Options 按实际负载调整。

自动 RECVACK 表示网络接收,不表示监听器已处理。即使没有监听器或监听器落后,也会确认接收;EasySDK 没有持久收件箱,需要可靠恢复时由业务层提供存储、去重和补偿。不要忽略 Lagged

发送不会自动重放。超时、取消或断线时服务端是否接受可能未知;业务决定重试时保留 SendOptions.client_msg_no 并遵循服务端幂等约定。SendOptions 还支持 Header、Setting 和 Topic;red_dot 默认 true,尊重显式 false。群聊使用 ChannelType::Group,成员与权限由后端准备。

群聊与成员管理

可信后端通过 Product HTTP 建群并管理成员。后端确认入群后,使用群 ID 和 ChannelType::Group 发送:

let ack = client.send(
    "project-team",
    wukong_easy_sdk::ChannelType::Group,
    serde_json::json!({"type": 1, "content": "大家好!"}),
).await?;

接收方仍通过 Event::Message 获取消息,检查 channel_typechannel_idfrom_uidpayload。SENDACK 表示服务端接受消息,不表示每位成员已经处理。 禁止陌生人发送的群中,非成员或已移除成员发送会得到 Error::Server { code: 3 }, 被加入黑名单的成员得到原因码 4。应用应处理这些权限错误;成员管理凭证由可信后端保管。

弱网、队列与退出

Backpressure 表示请求未获准进入发送队列:应用应限制生产并发,避免不断堆积重试任务。 发送超时可能发生在另一客户端已经收到消息之后,应保留“结果未知”,由应用依据幂等约定决定是否重试。 RecvError::Lagged 表示监听器错过了事件,需要应用侧补偿;自动 RECVACK 不提供持久化收件箱。

仓库验收工具反复执行两个客户端的 WSS 生命周期,注入每个数据块 20/40/60 ms 延迟、 返回方向暂停、双向黑洞和连接中断。通过两个待完成 SEND 检查准入上限,通过保留 16 个事件的慢监听器检查显式丢失通知。验收使用 800 ms SEND 超时和 3 秒 pong 超时, 这些是测试参数,与 SDK 默认值不同。

代理侧有界 WebSocket 审计要求每轮恰好 58 个出站 SEND,与 58 次已核对投递匹配, 避免服务端去重掩盖额外重发;审计仅保留计数,不支持的帧格式直接使验收失败。

每轮销毁两个客户端,并等待两个代理的连接流和任务都归零,再采样 Rust 探针的 RSS 和文件描述符。前三轮预热后,固定允许相对基线增加 64 MiB RSS 和八个文件描述符。 这些有限观测不证明生产容量、多日稳定性或不存在任何资源泄漏。

验证独立下载的正式包时,先按仓库 README 准备固定服务端提交,再运行:

python3 tests/acceptance/run.py --server-source test-server --distribution registry --seconds 120 --network-seconds 1800 --output .acceptance/network-1800s.json

弱网循环本身至少运行 30 分钟,构建和已有单聊/群聊检查另计。日常 CI 默认循环 30 秒且至少四轮,长程模式由手动运行明确选择。回执的 network 字段记录故障结果、 恢复耗时、资源采样和清理结果。资源观测支持 Linux 和 macOS。

WSS 私有 CA

使用私有 CA 时,将 DER 根证书字节加入 Options.additional_root_certificates。 最多 16 张,每张不超过 64 KiB,不接收私钥;公共 WebPKI 根仍保留,域名与有效期校验始终启用。

let options = wukong_easy_sdk::Options {
    additional_root_certificates: vec![std::fs::read("company-root.der")?],
    ..Default::default()
};

持续收发需要包含 WebSocket 缓冲区修复 的服务端版本。 单节点集群验收固定 27a39f15bf163b433f417b78ab6bfc6e589585e5;旧版在连续收发时可能损坏入队报文并断开连接。

协议与上线前检查

  • 发送使用字符串请求 ID、camelCase 字段与 Base64 UTF-8 JSON;接收支持 JSON 对象及 Base64 JSON;结果兼容 camelCase、snake_case 及同时出现的两种拼写。
  • SDK 默认静默;Auth 的 Debug 和 SDK 错误不包含 Token、URL、原始帧或服务端错误正文。完整事件包含业务数据,不应直接打印。
  • 自定义事件暴露 idevent_typetimestampdata;支持解析不意味着服务端部署一定产生某类事件。
  • 上线前检查 WSS 证书、代理 Upgrade、Token 撤销/轮换、目标操作系统、丢包、队列上限和退出;单次在线验收不代替长期稳定性和容量测试。

三节点集群验收

独立的集群验收脚本 固定包含跨节点成员缓存修复的服务端 f041174a042b4a96179218571e06c04bb64cf1ca,启动三个隔离进程, 配置 256 个 Hash Slot、12 个逻辑 Slot 和三个 Slot 副本。 四个 Rust 客户端通过私有 CA 验证的 WSS,分别接入节点 1、2、3、2。 验收覆盖前三个客户端之间全部六条有向单聊路径、群成员投递、频道隔离、 移除成员与黑名单拒绝,以及节点重启后的权限保持。

脚本先阻断发送者的返回流量,验证 SEND 超时而接收者已经收到消息; 随后强制终止接入节点 1,使用原地址、原数据目录重启。 原有客户端必须自动重连到节点 1。Slot 权威变更会清空易失的在线路由, 因此脚本还要求从三个 API 入口连续两个 25 秒心跳周期都能查询到四个用户在线, 再检查跨节点投递:观察窗口为 50 秒,整个路由检查最多等待 100 秒。 回执分别记录连接恢复时间和路由稳定检查完成时间。 独立的 WebSocket 出站计数必须与每次应用 SEND 匹配,包含被拒绝和结果未知的发送, 避免服务端去重掩盖客户端自动重发。应用不会重试结果未知的消息。试运行曾在 CONNECT 恢复后立即发送, 观察到接收者已不在在线路由列表、消息有 SENDACK 却没有投递。 因此 CONNECT 和 SENDACK 都不能作为端到端恢复或接收者已收到消息的证明。

这一限时场景验证原地址恢复后的重连,不提供其他地址自动选择, 也不代表节点停机期间连续投递、指定 Channel Leader 切换、网络分区、离线补拉或大群容量已经验证。 每个投递阶段至少观察被排除的客户端 500 ms。 测试配置为 SEND 超时 3 秒、连接超时 2 秒、PONG 超时 10 秒、最多重连 100 次、 重试退避 100–500 ms,与 SDK 默认值不同。

# 在 WuKongEasySDK-Rust 中运行;服务端使用下文记录的精确、干净源码。
RUSTUP_TOOLCHAIN=1.86.0 python3 tests/acceptance/cluster.py \
  --server-source ../test-server --distribution registry --seconds 600

CI 分别对源码和正式包运行 60 秒工作负载,手动运行可选择 600 秒。 每次运行还包括开始时的权限与故障检查,并在计时工作负载中再次终止、重启接入节点。 只有所有客户端完成销毁、查询四个用户得到空的在线状态列表,且所有自有服务端进程、 代理连接和任务完成清理,才算通过。回执分别记录正式包、验收脚本和服务端身份。

600 秒正式包三节点回执 分别记录干净验收工具提交 6b533a25ff0c61548a3f90dd36fa2562118f8f21、 精确的 0.1.0 正式包和上述已合并的服务端提交。 workload_seconds 包含中途故障与路由观察;fault_to_reconnected_ms 记录连接恢复时间,fault_to_stable_routes_ms 还包含 50 秒的路由观察窗口。 application_sends 必须等于 wire_sendsdeliveries 按所有预期接收者计数, 群聊多接收者投递会使其与发送次数不同。通过必须完成至少连续 600 秒工作负载、 两次故障阶段和全部清理,中断的候选版本运行不会拼接计入。现有 0.1.0 包不变。

验证记录

2026-09-08,crates.io wukong-easy-sdk 0.1.0 发布自源码 5b4a59cdbb66a9e0c3878e73ba4656f08ee05c6b。独立空缓存消费者从公共 registry 下载正式包,编译通过本页中英文的连接收发和私有 CA 示例。下载包 SHA-256 为 0029747f10b86f566e2d659535df0954114769a90962e562fb522a95e5508719,与 GitHub Release 附件一致。

随后完成了 crates.io 正式包端到端验收:macOS Rust 1.86 的独立空缓存消费者下载精确发布包,通过 Rust/Rust 收发、错误 Token 拒绝,以及 120 秒 Rust/JS WSS 验收,确认 1,747 次回显和三次断网恢复,无重复或事件丢失,所有进程完成清理。一次中断发送保留结果未知语义。正式包回执 分别记录验收工具源码 cfa48a038c2cfd56948ace43afe3b2f5f91dace3 与发布包源码;服务端仍为 27a39f15bf163b433f417b78ab6bfc6e589585e5

新增的 正式包群聊验收 同样通过:四个 Rust 客户端、两个群、十个场景, 确认 15 次预期投递。覆盖成员投递与群间隔离、非成员拒绝、成员添加/移除/重新加入、 黑名单拒绝与解除,以及四个客户端断网自动重连后的成员权限。每条消息都匹配 SENDACK 标识、发送人、群和正文,未观察到重复或事件丢失;每个场景至少观察 500 ms,检查被排除身份未收到消息。群聊回执 绑定干净验收工具提交 0262de9454603ee528dd2d9d9f236dec89e8df2a,使用相同的正式包 与服务端版本、macOS Rust 1.86;配套 120 秒 Rust/JS 验证确认 2,608 次回显、三次恢复, 一次中断发送保留结果未知语义。所有客户端和测试进程均完成清理。这是四客户端的 单节点集群验证,不证明大群容量、离线恢复或跨节点路由;本次无需修改 SDK 运行时代码或发布新包。

30 分钟弱网与资源回执 单独保留验收工具提交 2bd61a986f4f69418b3ec95de6c272408fc6f3b5,正式包仍为上述 0.1.0,服务端仍使用上述固定提交。network.cycles 记录出站 SEND 数、预期投递、 显式超时/背压/监听器落后及恢复结果,network.samples 记录连接清空后的 RSS 和文件描述符。 成功回执要求完整 1,800 秒循环、最后一轮完成和进程清理,不将多次短程运行拼接为长程结果。 资源阈值、测试参数和有限观测范围见上文。

同一源码通过 跨平台 CI:27 项协议、生命周期和 TLS 测试,以及 Clippy、API 文档和打包校验。覆盖 Linux Rust 1.86/stable、macOS 和 Windows stable。

真实服务端 CI 使用 WuKongIM 27a39f15bf163b433f417b78ab6bfc6e589585e5 的 256 Hash Slot 单节点集群,开启 Token 校验并拒绝错误 Token。Rust/Rust 验证通过后,与 npm easyjssdk@2.0.4 完成 120 秒 WSS 收发、3,012 次 Unicode 回显、三次断网恢复及清理,未观察到重复或事件丢失。五项 TLS 测试另行验证可信 CA、不可信 CA、错误域名、过期证书和无效配置。

另有源码 f30f1b32d0628f1e909fc21da704e5e49bc9f63e 的 600 秒 macOS 验收记录:13,758 次回显及三次恢复。该历史记录与正式包安装验证分别保留。旧服务端 132e46209d98fa0425cc0f88e7a97080cdad044d 仅通过早期短程验证,未通过持续收发,不应继承修复后服务端的验收结论。

公共 registry 正式包端到端验收与历史源码验证各自证明对应范围,不代表物理设备、离线恢复、容量或多日稳定性。返回 WuKongEasySDK 查看其他平台。

本页内容