@yidingdian/raft-cloud-sdk
v0.3.0-beta.3
Published
raft 中间件北向客户端 SDK:经 queue binding(appCmdQ/recvQ) 消费 raft 暴露的 Thing,封装读写/批量/订阅与 raft 私有 *Ex 扩展
Readme
@yidingdian/raft-cloud-sdk
上层业务服务与 raft 影子设备交互的 Client。任何 IoT 业务服务(或独立测试/调试工具)经它连接 raft,无需感知队列细节。
[业务服务 / 测试工具] ──(raft-cloud-sdk)──▶ appCmdQ ──▶ raft-cloud ──▶ 设备
◀────────────── recvQ ◀── (event/property/td 上报)实现机制(node-wot 与
*Ex旁路的分工等)见 DESIGN.md;使用者无需了解。
安装
npm install @yidingdian/raft-cloud-sdk需要宿主提供以下 peer 依赖:@yidingdian/core、@yidingdian/binding-queue、raft-common。
快速开始
const {RaftClient} = require('@yidingdian/raft-cloud-sdk');
const raft = new RaftClient({
redisOptions: { host, port, password },
redisDevDb: 3,
onProperty: ({sn, name, data, forms}) => { /* 上行属性上报 */ },
onEvent: ({sn, name, data, forms}) => { /* 上行事件上报 */ },
});
await raft.ready(); // 等 Servient/队列 binding 就绪
await raft.getBySn(sn); // 取影子 TD 并就绪(缓存优先)
await raft.readProperty(sn, 'staticInfo', {waitTTL: 5000});
await raft.readAllProperties(sn);
await raft.writeMultipleProperties(sn, null, {bright: 80, cct: 4000});
await raft.invokeAction(sn, 'blink', {times: 3});构造
new RaftClient(options)| 选项 | 默认 | 说明 |
|---|---|---|
| redisOptions | (必填) | ioredis 连接项 |
| redisDevDb | — | 设备/队列所在 redis db |
| appCmdQueueName | appCmdQ | 下行命令队列名 |
| recvQueueName | recvQ | 上行上报队列名 |
| skipRecvQueue | true | 宿主自带 recvQ worker 时跳过 binding 内建 worker |
| defaultReadWaitTtlMs | 30000 | 读默认 waitTTL |
| defaultJobWaitTtlMs | — | 写/动作的 caller 等待上限 |
| writeQueueTTLMs | 60000 | 陈旧写/动作在队列的丢弃阈值 |
| logger | no-op | 注入 bunyan 等 |
| bindingLogger | =logger | 队列 binding 库日志(常调更安静) |
| onProperty / onEvent | no-op | 上行上报钩子(或子类覆盖同名方法) |
| raftHttp | — | 影子库 HTTP 客户端 {host, port, accessKeyId, accessKeySecret};配了才有 listShadowThings / consumeAll,不配为纯队列模式 |
| deps | — | 注入 {Servient, QueueClientFactory, FairQueue}(测试/共享实例) |
API
就绪
| 方法 | 返回 | 说明 |
|---|---|---|
| ready() | Promise<this> | 等待 Servient/队列 binding 就绪;任何读写前须先 await(构造是异步初始化) |
TD 与设备句柄(含缓存管理)
RaftClient 内部缓存每台设备的影子 TD 与已建立的句柄,生命周期由宿主驱动:
- 填充:宿主收到上行
td通知时调cacheTd(sn, td)写入缓存; - 读取:
getBySn/getTd缓存优先,未命中才经队列getTd向 raft 拉取; - 失效:设备注销时调
invalidateTdCache(sn),清缓存并销毁句柄。
| 方法 | 返回 | 说明 |
|---|---|---|
| getBySn(sn) | 设备句柄 \| null | 取影子 TD 并就绪(缓存优先,未命中经队列 getTd);读写前须先调用 |
| getConsumedThing(td) | 设备句柄 | 已持有 TD 时直接 consume(getBySn 的下层;收到带全量 TD 的 td 通知时用) |
| getTd(sn) | td \| null | 取影子 TD(缓存优先) |
| getTd(sn, {full:true}) | td \| null | 取原始 TD(完整 schema)——仅供调试/查看:forms 是设备侧协议形态不可 consume,不进缓存,服务端多一次 DB 读 |
| cacheTd(sn, td) | — | 收到上行 td 通知时写入 TD 缓存 |
| destroyHandle(sn) | — | 销毁设备句柄(停订阅),保留 TD 缓存;用于 TD 更新时销毁旧句柄再以新 TD 重新 consume |
| invalidateTdCache(sn) | — | 设备注销:清 TD 缓存 + 销毁句柄 |
设备发现 / 批量预热(HTTP 影子库,需 raftHttp)
逐 SN 的队列 getTd 只能点查、需已知 SN 且依赖影子库命中;批量列表/预热走 raft 的 HTTP 影子接口
(/api/v1/shadow/things,accessKey 签名)。这是业务侧启动 consumeThings 的通用等价物。
| 方法 | 返回 | 说明 |
|---|---|---|
| listShadowThings({page, perpage, sns, thingIds}) | {count, rows} | 分页/条件列影子(rows[].td + rows[].m_shadow_device.{sn,status}),设备发现用 |
| consumeAll({perpage, batchSize, onProgress}) | {count} | 分页拉全量 → 缓存 TD + consume;之后 getBySn/读写无需逐 SN getTd |
| thingInfo(sn) | {tdvr, productId} \| null | 已就绪设备的型号/版本 |
| getRepeatObserve(sn) | Map \| null | 设备的「属性 → 重复上报频率」映射(供重复 observe 处理) |
读属性
| 方法 | 返回 |
|---|---|
| readProperty(sn, name, {waitTTL}) | value |
| readMultipleProperties(sn, names, {waitTTL}) | {name: value} |
| readAllProperties(sn, {waitTTL}) | {name: value} |
写属性 / 调用动作
| 方法 | 说明 |
|---|---|
| writeProperty(sn, name, value) | 单属性写 |
| writeMultipleProperties(sn, null, values) | 批量写;values 为 {name: value}(第二参当前忽略,传 null) |
| writeAllProperties(sn, values) | 同上,写全部 |
| invokeAction(sn, action, input) | 返回动作输出(无则 undefined) |
上行上报
设备上报(属性变更 / 事件)经 onProperty / onEvent 钩子回调,报文结构:
{ sn, name, data, forms }两种提供方式:构造时传 onProperty / onEvent,或子类覆盖同名方法。
行为契约
waitTTL快速失败:读的options.waitTTL是 caller 愿等待毫秒数;预计超时立即失败返回,不阻塞。缺省defaultReadWaitTtlMs。readAllProperties依赖设备能力:设备支持"一条指令返回全部"时走单指令;否则自动退回逐属性读取。- 并发读自动合并:同一属性的并发读会合并为一次设备读(各 caller 的
waitTTL仍各自独立)。 - 单属性优化:
writeMultipleProperties/writeAllProperties传入仅一个属性时,等价于writeProperty。
上行投递确认(status)
中间件稳态只在状态变化时投递 status,一条丢了就没有第二次。宿主消费掉一条就回记, 中间件在对账窗口到点时比对:追平就跳过,落后才重推。不回记 = 每个窗口被重推一次。
const {createAckWriter} = require('@yidingdian/raft-cloud-sdk');
const ack = createAckWriter({redisOptions, redisDevDb, logger});
await ack.ackStatus(sn, deviceTs); // deviceTs = 上行 job 的 timestamp用 SDK 内建 recvQ worker 的宿主无需自己调;自带 worker 的宿主要在两处回记:
- 该条 status 处理成功之后;
- 因为手里已有更新的状态而丢弃它时——对中间件而言它同样已送达,漏掉会导致永久重推。
约定:
- 只进不退,回记比现值旧的 ts 会被忽略;
- 回记失败只告警,最坏是被多推一次,不该影响消息处理;
- 键名与存放位置是 SDK 的实现细节,宿主不要自己拼 key。
上行拥塞状况(下发前问一句)
有些下发会引发一大批设备同时上报(典型:网关上线后下发它的子设备列表,子设备随即全量上线)。 上行正堵的时候再灌进去,只会把堵的时间拉长。中间件每秒产出一份拥塞快照,宿主读它决定发还是延后。
const {createUplinkStatusReader} = require('@yidingdian/raft-cloud-sdk');
const uplink = createUplinkStatusReader({redisOptions, redisDevDb, logger});
const {defer, retryAfterMs, level, reasons} = await uplink.check();
if (defer) scheduleRetry(retryAfterMs); // level: busy / overload,reasons 指出是谁在堵
else await send();get() 返回原始快照,自己定阈值时用:
| 字段 | 含义 |
| --- | --- |
| level | idle / busy / overload / unknown(拿不到快照) |
| busy / retryAfterMs | 判定结果与建议重试间隔 |
| reasons | 触发判定的指标名,如 ['gateQueued','elu'] |
| pendingDevices | 上行队列里还有多少台设备的上报没被取走;null = 未知。默认不参与判定,见下 |
| gateQueued | 中间件上行准入闸的排队数 |
| ratePerSec / elu | 上行速率与事件循环利用率,趋势参考 |
约定:
pendingDevices默认只上报、不参与level判定:它量的是对端还有多少没取走,稳态基线随负载浮动, 绝对阈值贴上基线就会把调用方锁死。要判它得换算成"还要多久排空"= 深度 ÷ 自己的消费速率—— 分母只有消费端知道(中间件量到的是入队速率),所以这一步留给调用方,get()里取原始深度自己算;- 拿不到快照(中间件未升级/停摆)默认不延后,读取器只告警一次;要反过来(宁可不发)传
deferWhenUnknown: true; - 快照过老(默认 15s)按未知处理,不拿旧数当"不堵";
- 同一瞬间多次查询共用一次 redis 读(默认 500ms 缓存);
- 键名与存放位置是实现细节,宿主不要自己拼 key。
业务子类扩展
业务把落库 / 固件规避等留在子类,SDK 主干保持通用:
class MyConsumer extends RaftClient {
_normalizeWriteValue(sn, thing, name, value) { /* 产品/固件特例 */ return value; }
async onProperty(report) { /* 转交业务上行处理管线 */ }
async onEvent(report) { /* ... */ }
}测试
npm test # 冒烟:零业务耦合校验 + 注入 fake 依赖验证命令编码(不连 redis)