Skip to content

3.13 TopicPubSub 订阅发布类

概述

TopicPubSub 提供基于 WebSocket 的实时订阅发布能力。 推荐入口是 Arm::topicPubSub

  1. 先调用 Arm::Connect()
  2. 再调用 arm.topicPubSub.Connect()
  3. 启动接收
  4. 发起状态 / 寄存器 / IO 订阅

推荐用法

cpp
Arm arm;  // 创建机器人会话对象
STATUS_CODE ret = arm.Connect("10.27.1.254");  // 先连接控制器或路由地址
if (ret != STATUS_CODE::OK) {  // 判断 Arm 连接是否失败
    return;  // 连接失败时退出当前业务流程
}  // 结束 Arm 连接结果判断
ret = arm.topicPubSub.Connect();  // 在 Arm 连接成功后建立 SubPub WebSocket 连接
if (ret != STATUS_CODE::OK) {  // 判断 SubPub 连接是否失败
    arm.Disconnect();  // SubPub 连接失败时断开 Arm 会话
    return;  // 退出当前业务流程
}  // 结束 SubPub 连接结果判断
arm.topicPubSub.StartReceiving(  // 启动后台接收并注册消息回调
    [](const Json::Value& message) {  // 定义收到 JSON 消息时执行的回调
        // 处理消息  // 在这里解析 message,耗时逻辑应转交业务线程
    }  // 结束消息回调
);  // 结束接收启动调用
arm.topicPubSub.SubscribeStatus({RobotTopicType::JOINT_POSITION});  // 订阅机器人关节位置主题

场景化示例

下面几组片段按 “连接接收、订阅主题、主动收发、清理连接” 交叉覆盖本页 API。片段默认承接已经 Arm::Connect() 成功的 arm 对象。

建立连接并启动回调接收

cpp
STATUS_CODE connectRet = arm.topicPubSub.Connect("", 10);  // 使用 Arm 已绑定的地址建立 SubPub 连接,握手超时 10 秒
if (connectRet != STATUS_CODE::OK || !arm.topicPubSub.IsConnected()) {  // 判断 SubPub 是否连接成功
    return;  // 连接失败时退出当前流程
}  // 结束 SubPub 连接判断
STATUS_CODE receiveRet = arm.topicPubSub.StartReceiving(  // 启动后台接收并安装 JSON 消息回调
    [](const Json::Value& message) {  // 定义收到消息时的处理函数
        const Json::Value copy = message;  // 示例中只复制消息,实际业务可转交自己的队列
        (void)copy;  // 避免示例变量未使用
    }  // 结束消息回调
);  // 结束后台接收启动调用

订阅状态、寄存器和 IO

cpp
STATUS_CODE statusRet = arm.topicPubSub.SubscribeStatus(  // 订阅机器人状态主题
    {RobotTopicType::JOINT_POSITION},  // 订阅关节位置主题
    200  // 使用 200Hz 订阅频率
);  // 结束状态订阅调用
STATUS_CODE regRet = arm.topicPubSub.SubscribeRegister(  // 订阅寄存器主题
    RegTopicType::R,  // 订阅 R 数值寄存器
    {1, 2},  // 订阅 1 号和 2 号寄存器
    100  // 使用 100Hz 订阅频率
);  // 结束寄存器订阅调用
STATUS_CODE ioRet = arm.topicPubSub.SubscribeIo(  // 订阅 IO 主题
    {{IOTopicType::DI, 1}},  // 订阅 1 号数字输入
    50  // 使用 50Hz 订阅频率
);  // 结束 IO 订阅调用

主动发送文本并从队列取消息

cpp
STATUS_CODE sendRet = arm.topicPubSub.SendText("{\"cmd\":\"ping\"}");  // 发送原始文本,通常只在调试或自定义协议时使用
auto [message, messageRet] = arm.topicPubSub.Receive(1000);  // 从 SDK 接收队列取一条消息,最多等待 1000ms
(void)sendRet;  // 示例中保留发送结果,实际代码应检查错误码
(void)message;  // 示例中保留收到的消息,实际代码应解析 JSON 内容
(void)messageRet;  // 示例中保留接收结果,实际代码应区分超时和断连

移除回调并断开连接

cpp
STATUS_CODE removeRet = arm.topicPubSub.RemoveMessageHandler();  // 移除当前消息回调,保留接收队列能力
STATUS_CODE disconnectRet = arm.topicPubSub.Disconnect();  // 断开 SubPub WebSocket 连接
(void)removeRet;  // 示例中保留移除回调结果,实际代码应检查错误码
(void)disconnectRet;  // 示例中保留断开结果,实际代码应检查错误码

示例代码

cpp17/sub_pub_basic/src/main.cpp
cpp
#include "multi_instance_isolation/run.h"
#include "subscribe_topics/run.h"
#include "send_receive_text/run.h"

int main(void)
{
    // [ZH] 默认只调用一个门面方法;如需体验其他接口,请把下一行替换成下面任意一行。
    // [EN] The main function calls only one facade by default. Replace the next line with any line below to try other APIs.
    return RunSubPubBasicSubscribeTopics();
    // return RunSubPubBasicSendReceiveText();
    // return RunSubPubBasicMultiInstanceIsolation();
}

接口总览

cpp
Connect(const std::string& teachPanelIp = "", int32_t timeoutSecs = 10) -> STATUS_CODE
IsConnected() const -> bool
StartReceiving(const MessageHandler& handler) -> STATUS_CODE
SubscribeStatus(const std::vector<ROBOT_TOPIC_TYPE>& topics, int32_t frequency = 200) -> STATUS_CODE
SubscribeRegister(const REG_TOPIC_TYPE& regType, const std::vector<int32_t>& regIds, int32_t frequency = 200) -> STATUS_CODE
SubscribeIo(const std::vector<std::pair<IO_TOPIC_TYPE, int32_t>>& ioList, int32_t frequency = 200) -> STATUS_CODE
SendText(const std::string& text) -> STATUS_CODE
RemoveMessageHandler() -> STATUS_CODE
Receive(int32_t timeoutMs = 5000) -> std::pair<Json::Value, STATUS_CODE>
Disconnect() -> STATUS_CODE
方法输入输出关键行为
Connect可选地址、握手超时秒数STATUS_CODE空地址时优先使用 Arm::Connect() 期间绑定的 TP/IP
IsConnectedbool只反映 sub_pub WebSocket 连接状态
StartReceiving单个 JSON 回调STATUS_CODE启动后台接收;后一次回调会覆盖前一次
SubscribeStatustopic 列表、频率STATUS_CODE发送 addTopic 命令
SubscribeRegister寄存器类型、ID 列表、频率STATUS_CODE发送 addRegTopic 命令
SubscribeIoIO 类型 + 编号列表、频率STATUS_CODE发送 addIoTopic 命令
SendText原始文本STATUS_CODE发送原始文本
RemoveMessageHandlerSTATUS_CODE移除当前活动回调,不影响 SDK 接收队列
Receive超时毫秒std::pair<Json::Value, STATUS_CODE>从 SDK 接收队列取下一条消息
DisconnectSTATUS_CODE关闭当前 WebSocket 连接

详细语义

Connect

签名

cpp
STATUS_CODE Connect(const std::string& teachPanelIp = "", int32_t timeoutSecs = 10);
说明
teachPanelIp显式传入时直接使用;空字符串时,优先使用 Arm::Connect() 已绑定的目标地址
timeoutSecsWebSocket 握手超时,默认 10 秒
返回OK / INVALID_IP_ADDRESS / 连接阶段其他错误码

限制与行为

  • 如果你是通过 arm.topicPubSub 使用,通常直接调用无参 Connect() 即可。
  • 如果你单独实例化 TopicPubSub ,则必须显式传入地址。
  • 固定使用代理 WebSocket 端口 5609
  • 如果当前已经连上同一 TopicPubSub 实例,再次调用直接返回 OK ,沿用当前连接。
  • 多个 Arm / TopicPubSub 实例可以在同一进程内并发共存;断开其中一个实例时,其他实例继续使用各自的 WebSocket 网络环境。

多实例隔离示例:

cpp17/sub_pub_basic/src/multi_instance_isolation/run.cpp
cpp
#include <iostream>
#include <json/json.h>
#include <string>
#include <utility>

#include "arm_api.h"
#include "status_code.h"

#include "multi_instance_isolation/run.h"

namespace {

std::pair<Json::Value, STATUS_CODE> SendPingAndReceive(
    Arm& arm,
    const std::string& source
)
{
    Json::Value request;
    request["cmd"] = "ping";
    request["param"]["source"] = source;

    Json::StreamWriterBuilder builder;
    builder["indentation"] = "";
    STATUS_CODE sendRet = arm.topicPubSub.SendText(Json::writeString(builder, request));
    if (sendRet != STATUS_CODE::OK) {
        return std::make_pair(Json::Value(), sendRet);
    }

    return arm.topicPubSub.Receive(5000);
}

} // namespace

/**
 * 多实例 TopicPubSub 隔离门面。
 * @return 0 表示成功,否则返回 1。
 */
int RunSubPubBasicMultiInstanceIsolation(void)
{
    // [ZH] 请把下面地址替换成当前机器人控制器 IP 和示教器 IP。
    // [EN] Replace the following addresses with the current controller IP and teach panel IP.
    const std::string controllerIp = "10.27.1.2";
    const std::string teachPanelIp = "10.27.1.102";

    Arm armA;
    Arm armB;
    STATUS_CODE armConnectRetA = armA.Connect(controllerIp, teachPanelIp);
    STATUS_CODE armConnectRetB = armB.Connect(controllerIp, teachPanelIp);
    if (armConnectRetA != STATUS_CODE::OK || armConnectRetB != STATUS_CODE::OK) {
        std::cerr << "[cpp17_sub_pub] 多实例连接失败 / Multi-instance connect failed, A="
                  << static_cast<int>(armConnectRetA) << ", B="
                  << static_cast<int>(armConnectRetB) << "\n";
        armA.Disconnect();
        armB.Disconnect();
        return 1;
    }

    STATUS_CODE subPubConnectRetA = armA.topicPubSub.Connect(teachPanelIp, 10);
    STATUS_CODE subPubConnectRetB = armB.topicPubSub.Connect(teachPanelIp, 10);
    auto [ackA, ackRetA] = SendPingAndReceive(armA, "cpp17_multi_instance_A");
    STATUS_CODE disconnectRetA = armA.topicPubSub.Disconnect();
    auto [ackB, ackRetB] = SendPingAndReceive(
        armB,
        "cpp17_multi_instance_B_after_A_disconnect"
    );
    STATUS_CODE disconnectRetB = armB.topicPubSub.Disconnect();

    Json::StreamWriterBuilder builder;
    builder["indentation"] = "";
    std::cout << "[cpp17_sub_pub] ArmA Connect 状态码 / ArmA connect status code: "
              << static_cast<int>(subPubConnectRetA) << "\n";
    std::cout << "[cpp17_sub_pub] ArmB Connect 状态码 / ArmB connect status code: "
              << static_cast<int>(subPubConnectRetB) << "\n";
    std::cout << "[cpp17_sub_pub] ArmA Receive 状态码 / ArmA receive status code: "
              << static_cast<int>(ackRetA) << ", 消息 / Message: "
              << Json::writeString(builder, ackA) << "\n";
    std::cout << "[cpp17_sub_pub] ArmA Disconnect 状态码 / ArmA disconnect status code: "
              << static_cast<int>(disconnectRetA) << "\n";
    std::cout << "[cpp17_sub_pub] ArmB Receive 状态码 / ArmB receive status code: "
              << static_cast<int>(ackRetB) << ", 消息 / Message: "
              << Json::writeString(builder, ackB) << "\n";
    std::cout << "[cpp17_sub_pub] ArmB Disconnect 状态码 / ArmB disconnect status code: "
              << static_cast<int>(disconnectRetB) << "\n";

    armA.Disconnect();
    armB.Disconnect();
    std::cout << "[cpp17_sub_pub] 多实例隔离示例结束 / Multi-instance isolation example finished\n";

    const bool success =
        subPubConnectRetA == STATUS_CODE::OK &&
        subPubConnectRetB == STATUS_CODE::OK &&
        ackRetA == STATUS_CODE::OK &&
        disconnectRetA == STATUS_CODE::OK &&
        ackRetB == STATUS_CODE::OK &&
        disconnectRetB == STATUS_CODE::OK;
    return success ? 0 : 1;
}

IsConnected

签名

cpp
bool IsConnected() const;

限制与行为

  • 这里检查的是 sub_pub WebSocket,不是 Arm 的 HTTP 会话。
  • Arm::Connect() 成功后, arm.topicPubSub.IsConnected() 仍然是 false ;只有 arm.topicPubSub.Connect() 成功后才会变成 true

StartReceiving

签名

cpp
STATUS_CODE StartReceiving(const MessageHandler& handler);

输入

参数类型说明
handlerstd::function<void(const Json::Value&)>消息回调

输出

  • 成功:返回 OK
  • 未连接:返回 NOT_CONNECTED

限制与行为

  • 同一个 TopicPubSub 实例只保留一个活动回调,后一次调用会覆盖前一次回调。
  • 回调收到的是已经解析好的 Json::Value
  • 回调在 WebSocket 消息到达时触发。
  • SDK 接收队列上限是 100 条,超过上限时会丢弃最旧消息。
  • 回调应尽量轻量,耗时逻辑请转交业务线程。
  • 不要在回调里直接调用 Arm::Connect()Arm::Disconnect()TopicPubSub::Connect()TopicPubSub::Disconnect()StartReceiving()RemoveMessageHandler() ;这类重入操作可能返回 OTHER_ERR 或不生效。

SubscribeStatus

签名

cpp
STATUS_CODE SubscribeStatus(
    const std::vector<ROBOT_TOPIC_TYPE>& topics,
    int32_t frequency = 200
);
说明
topics要订阅的主题列表,常量见 ROBOT_TOPIC_TYPE
frequency订阅频率,单位 Hz,默认 200
订阅行为添加机器人状态主题

SubscribeRegister

签名

cpp
STATUS_CODE SubscribeRegister(
    const REG_TOPIC_TYPE& regType,
    const std::vector<int32_t>& regIds,
    int32_t frequency = 200
);
说明
regType寄存器类型,见 REG_TOPIC_TYPE
regIds寄存器编号列表
frequency订阅频率,单位 Hz
订阅行为添加寄存器主题

SubscribeIo

签名

cpp
STATUS_CODE SubscribeIo(
    const std::vector<std::pair<IO_TOPIC_TYPE, int32_t>>& ioList,
    int32_t frequency = 200
);
说明
ioList每项是 (ioType, ioId)
frequency订阅频率,单位 Hz
订阅行为添加 IO 主题

SendText

签名

cpp
STATUS_CODE SendText(const std::string& text);

用于直接发送原始文本。 如果你只做状态、寄存器或 IO 订阅,通常优先使用上面的订阅接口。

RemoveMessageHandler

签名

cpp
STATUS_CODE RemoveMessageHandler();

限制与行为

  • 同一个 TopicPubSub 实例只有一个活动回调,所以这里移除的是 “当前回调”。
  • 移除回调后, Receive() 仍然可以继续从 SDK 接收队列取消息。
  • 即使当前未连接,也会返回 OK

Receive

签名

cpp
std::pair<Json::Value, STATUS_CODE> Receive(int32_t timeoutMs = 5000);
说明
timeoutMs超时时间,默认 5000ms
成功返回Json::Value + OK
超时返回Json::Value + SUB_PUB_RECEIVE_TIMEOUT
断连返回Json::Value + NOT_CONNECTED

Disconnect

签名

cpp
STATUS_CODE Disconnect();

断开当前 WebSocket 连接。 如果当前尚未连接,返回 OK