3.13 TopicPubSub 订阅发布类
概述
TopicPubSub 提供基于 WebSocket 的实时订阅发布能力。 推荐入口是 Arm::topicPubSub :
- 先调用
Arm::Connect() - 再调用
arm.topicPubSub.Connect() - 启动接收
- 发起状态 / 寄存器 / 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; // 示例中保留断开结果,实际代码应检查错误码示例代码
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 |
IsConnected | 无 | bool | 只反映 sub_pub WebSocket 连接状态 |
StartReceiving | 单个 JSON 回调 | STATUS_CODE | 启动后台接收;后一次回调会覆盖前一次 |
SubscribeStatus | topic 列表、频率 | STATUS_CODE | 发送 addTopic 命令 |
SubscribeRegister | 寄存器类型、ID 列表、频率 | STATUS_CODE | 发送 addRegTopic 命令 |
SubscribeIo | IO 类型 + 编号列表、频率 | STATUS_CODE | 发送 addIoTopic 命令 |
SendText | 原始文本 | STATUS_CODE | 发送原始文本 |
RemoveMessageHandler | 无 | STATUS_CODE | 移除当前活动回调,不影响 SDK 接收队列 |
Receive | 超时毫秒 | std::pair<Json::Value, STATUS_CODE> | 从 SDK 接收队列取下一条消息 |
Disconnect | 无 | STATUS_CODE | 关闭当前 WebSocket 连接 |
详细语义
Connect
签名
cpp
STATUS_CODE Connect(const std::string& teachPanelIp = "", int32_t timeoutSecs = 10);| 项 | 说明 |
|---|---|
teachPanelIp | 显式传入时直接使用;空字符串时,优先使用 Arm::Connect() 已绑定的目标地址 |
timeoutSecs | WebSocket 握手超时,默认 10 秒 |
| 返回 | OK / INVALID_IP_ADDRESS / 连接阶段其他错误码 |
限制与行为
- 如果你是通过
arm.topicPubSub使用,通常直接调用无参Connect()即可。 - 如果你单独实例化
TopicPubSub,则必须显式传入地址。 - 固定使用代理 WebSocket 端口
5609。 - 如果当前已经连上同一
TopicPubSub实例,再次调用直接返回OK,沿用当前连接。 - 多个
Arm/TopicPubSub实例可以在同一进程内并发共存;断开其中一个实例时,其他实例继续使用各自的 WebSocket 网络环境。
多实例隔离示例:
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);输入
| 参数 | 类型 | 说明 |
|---|---|---|
handler | std::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 。