资源与协议接入手册
0. 什么是「资源」
资源是这个平台里唯一的一类实体:一个有 ID、有状态、会上报数据、可以被单独授权和吊销的东西。
平台不再建立另一套实体分类。下列对象都可以直接建模为资源:
| 当作资源 | 它上报的是 | | --- | --- | | 一台传感器、电表、车机 | 温度、电量、位置 | | 一条产线、一个工位 | 节拍、良率、停机事件 | | 一个服务实例、一个爬虫、一个定时任务 | 心跳、处理量、错误 | | 一个业务对象(订单、工单、库位) | 状态变更 | | 一个仿真器、一段回放 | 和真实来源同形的数据 |
平台对这些一视同仁,因为它对资源"是什么"根本不做判断:
- 路由层不认识资源。 主题完全是你的,消息体完全是你的。
- 身份来自凭据,不来自载荷。
publisherId是resource:sensor-01还是resource:order-4471,对网关是同一件事。 - 归属可以从数据里解析。 接外部 broker 时,资源 ID 是按你写的规则从主题或 payload 里取出来的——那条规则里没有任何地方假设对面是硬件。
只要一个对象需要独立 ID、状态或数据流,就可以把它建模为资源。
1. 统一数据模型
所有协议最终写入同一个可靠通道:
space/{domain}/{space}/pub/{publisherId}domain 是租户或组,space 是空间(默认 -),publisherId 是发布者的权威身份(resource:sensor-01、tenant:acme),由凭据决定而不是由载荷声明。资源通过 REST 上报时落在 space/{租户}/-/pub/resource:{资源}。
每条事件包含 eventId、严格递增的 sequence、type、data 和服务端 createdAt。同一资源重发相同 eventId 时返回原事件,不会新增记录。遥测 data 必须为 JSON 对象,默认最大 1 MiB。
sequence 用于排序和游标,保证在同一资源通道内严格递增,但在多网关实例部署下不保证连续。消费端只应比较大小、记录“最后成功处理的 sequence”,不要假设下一条一定是 sequence + 1,也不要用它统计条数。
推荐资源生成 UUIDv4 作为 eventId,在收到协议确认前一直重试同一个 ID。成功后才能生成下一个 ID。
2. 资源注册和认证
2.0 先说清楚:资源密钥是可选的
路由层不认识资源。你拿租户密钥 etk_… 就能收发消息,一个资源都不注册也完全能用——这是这套架构本来的样子,不是退化模式。
那把租户密钥直接放进资源行不行?行,风险你自己权衡,说清楚是:
- 撬开任意一台,拿到的是整个租户的权限,不只是那一台的数据;
- 想吊销就只能换租户密钥,整队一起停;
- 出事后分不清是哪台发的——
publisherId只会是tenant:acme。
资源密钥 erk_… 解决的就是这三条,它是一个更窄的凭据,不是一道必须过的门。评估阶段、内部系统、资源本来就在可信网络里的场景,不用它是合理选择。
下面 2.1 是手工发证,2.2 是批量发证(引导凭据 / 一机一密)。
2.1 先有账号
登录、注册和会话由 @ptner/worker-auth 提供——和其他服务共用同一套,所以任何一处修好了,所有地方一起修好,而不是各有一份手写的。它的 storage-lib 适配器让它跑在这个网关配的任何后端上(SQLite / MongoDB / D1 / Bunny)。
三种方式通向同一个账号:
| 方式 | 需要什么 | | --- | --- | | Google | 配 GOOGLE_CLIENT_ID + GOOGLE_CLIENT_SECRET | | 邮件链接 | 配邮件发送;没配时链接直接回显(本地开发用) | | 邮箱 + 密码 | 无需额外配置 |
| 路径 | 用途 | | --- | --- | | /signup /login | 注册 / 登录 | | /tenant | 选租户 ID(账号建好之后那一步) | | /console | 控制台。只能用账号进——API 密钥是给程序的 |
账号和租户是两步。 注册只建账号,租户 ID 要另外选——它会出现在 MQTT 用户名和主题里,是你的选择,而且在一次 OAuth 跳转或一封邮件链接里没地方填。所以「有账号但还没有租户」是一个正常状态,它能跨浏览器重启存活。
选完租户不会发密钥。 控制台用你的账号会话,本来就不需要密钥;要密钥是因为你要跑一个程序,那就去「凭据」页建一把,名字和范围当场划清楚。控制台也不接受粘贴密钥登录:一把密钥贴进浏览器就是一个页面既无法限定范围也无法作废的长期凭据,而且事后看不出是谁在操作。丢了不用「换」——新建一把、部署、删掉旧的,中间两把并存,没有断档。用 Google 或邮件链接进来的账号可以在「账号」页补设一个密码,多一条登录路。
会话 cookie 是 HttpOnly + SameSite=Lax、没有 Domain(兄弟子域收不到),TLS 上才加 Secure。只存 token 的 SHA-256,拷走数据库拿不到活会话。会话解析出来的就是租户密钥同一个身份,所以所有接口对两种方式行为一致;同时带了 Authorization 头时显式密钥优先——否则已登录的浏览器就没法以管理员身份操作。
SameSite=Lax 仍允许跨站 GET 导航,本身不足以防 CSRF;所有改动操作都在 POST/PUT/DELETE 上,且要求 content-type: application/json(跨站表单发不出这个类型)。
管理员由部署方在 ADMIN_EMAILS 里列出真实邮箱(逗号分隔),那个人用自己的账号登录就是管理员,登录方式和其他人完全一样。这样做而不是给一个保留域的角色账号,是因为:保留域收不到邮件,那个账号只能配一个静态密码放在环境变量里——正好是账号体系要消除的共享密钥;而且角色账号让审计日志里每条管理操作都记在同一个名下,看不出是谁。名单是每次请求现查的,所以从里面删掉一个人立即生效,不用等他的会话过期。管理员没有自己的租户,登录后在「租户」页选一个。
ADMIN_API_KEY 保留,但只用于程序:脚本、CI、/metrics。
部署方可以 SIGNUP_ENABLED=false 关闭新账号注册(已有账号照样能登录,包括用邮件链接)。新账号必须通过邮件链接或 Google 验证邮箱;验证之后可以在「账号」页设置密码。PUBLIC_BASE_URL 决定邮件链接和 Google 回调指向哪里——不设就用实际绑定的地址并打一条警告,代理后面必须设。Google 回调要在 Google Cloud Console 里登记为 PUBLIC_BASE_URL + /v1/auth/google/callback,完全一致。
已经有自己的 Broker:不用动资源
资源已经连在你自己的 MQTT broker 上,为了用地球视图和审计而把上万台重新指向这里是不现实的。所以平台主动连过去:控制台「外部 Broker」页填地址、凭据、要订阅的主题,剩下的它自己做。
资源是从消息里认出来的,不靠凭据。 两种规则:
| 规则 | 写法 | 例子 | | --- | --- | --- | | 主题里的一级 | plant/{resourceId}/state | 主题 plant/boiler-7/state → boiler-7 | | payload 里的字段 | meta.dev | {"meta":{"dev":"boiler-7"}} → boiler-7 |
占位符也可以在一级里面(dev-{resourceId})。规则对不上的主题跳过,不会猜——跳过的条数在列表里看得到。
这一点是有意的:一机一密凭据是可选项,不是前提。接入连接本身由租户凭据认证,那是信任边界;资源 ID 是这条可信流里的归属信息,从数据里解析出来。没见过的资源第一次上报时自动登记。
地址、凭据、主题、规则都能改,在列表里点「修改」——同一个表单载入原值,保存后运行时会拆掉旧连接按新配置重连。密码框留空表示不改(API 从不把它读回来,所以也填不回去);要清掉就提交一个空密码。
几个实现上的事实:
- 消息里带
ts/timestamp/time就用它当观测时间(秒和毫秒都认),否则用到达时间。 - 保留消息(retained)按内容去重,但重放次数单独计,而且循环本身会报警。
broker 在每一次订阅时都会把保留消息完整重发一遍,而接入每次重连都会重新订阅。这种重放不是新读数——实测确认:实时发布给已订阅客户端时投递的 retain=false,只有新订阅从保留库拉取才是 true,所以按内容去重不会吞掉真实读数(资源周期性重发 retained 状态是 retain=false,照常入库)。
但去重会让「为什么一直在重放」这件事消失,那才是真问题。所以:重放次数单独计(列表里的「保留重放」,不和「重复」混在一起),而且重连循环本身被直接检测——窗口内超过阈值就把状态标成「反复重连」,并写一条 connector.flapping 审计(每窗口一条,不刷屏)。计数器随进程消失,审计不会。
阈值:CONNECTORS_FLAP_THRESHOLD(默认 5)/ CONNECTORS_FLAP_WINDOW_MS(默认 60 秒)。
- 每个客户端的 clientId 都不一样。MQTT broker 会把同 id 的旧连接踢掉,被踢的那个重连又把新的踢掉——一个 5 秒一轮的死循环,而每轮重连都重新拉一遍保留消息。会话是
clean的,固定 id 没有任何好处。 - broker 密码用
CREDENTIAL_PEPPER加密存储(和发证种子同一套),API 不会把它读回来。 - 状态由服务端推(WebSocket 的
status消息),不是轮询,也不用点刷新。「连接中」后面带着已经连了多久和第几次尝试——不然一个丢包不回的地址会连上两分钟,看起来和卡死一样。 - 不做集群:每个实例都跑就会重复收。
CONNECTORS_ENABLED指定唯一拨号的那个实例,默认开——单实例是常见部署。
浏览器为什么要票据,别的客户端不要
只有浏览器需要。 程序直接拿 etk_ 连就行——HTTP 放 Authorization 头,MQTT 放 CONNECT 的 password 字段,WebSocket 也放 Authorization 头(网关的升级处理里没带票据就去读头)。都不需要换票。
浏览器卡在一个地方:new WebSocket(url, protocols) 没有填 header 的参数——这是浏览器 API 的限制,不是 WebSocket 协议的。而且控制台是账号登录,页面手上只有一个 HttpOnly cookie,JS 读都读不到。
所以票据 ewt_ 的作用只有一个:把一个 JS 碰不到的凭据,换成一个 JS 能填进去的凭据。手上已经有 etk_ 的,不要换,直接用。
为什么不干脆让浏览器带 cookie 去连:cookie 按主机划分、不看端口,所以控制台的 cookie 本来就能到 MQTT 的端口;而 WebSocket 升级不受 CORS 约束,任何网站都能开一条连接,浏览器照样附上 cookie。那是真的洞,票据没有这个问题(跨站签不出来)。
票据一次性、默认 60 秒(WS_TICKET_TTL_MS)。第二次使用会被记进审计(ticket.replayed)——一次性买到的不是「重放失败」,是「有人知道出现了第二份拷贝」。
密码找回没有实现(要邮件模板和额外流程)——不过邮件链接登录本身就是可用的替代:能收到邮件就能进去,进去后在「账号」页改密码。
2.2 手工注册一个资源
在控制台「资源」页填一张表,页面直接给出资源密钥和填好参数的接入示例。开通租户、metadata 约定和界面说明见 控制台与自助接入。
等价的 API 调用,租户密钥和管理员密钥都可以:
curl -sS -X POST https://resources.example.com/v1/resources \
-H "Authorization: Bearer $TENANT_API_KEY" \
-H 'X-Tenant-Id: acme' \
-H 'Content-Type: application/json' \
-d '{"tenantId":"acme","resourceId":"sensor-01","name":"Warehouse sensor","metadata":{"model":"T100"}}'返回:
{
"resource": { "tenantId": "acme", "resourceId": "sensor-01", "status": "registered" },
"apiKey": "erk_..."
}再次注册相同资源会轮换密钥:旧密钥无法再发起新认证,用它建立的 MQTT / WebSocket / WebTransport 会话也会被服务端主动断开,资源应当以新密钥重连。服务端只保存带 pepper 的 SHA-256 结果,不保存明文。资源凭据的 MQTT/WebTransport 用户名统一为 tenantId/resourceId。
凭据的识别方式和权限边界:
| 角色 | 密钥 | HTTP 请求头 | MQTT username | 范围 | | --- | --- | --- | --- | --- | | 账号会话 | 登录后的 cookie | Cookie(浏览器自动带) | — | 等同于本租户的租户密钥 | | 管理员 | ADMIN_API_KEY | Authorization | — | 全部租户;/metrics、租户管理 | | 租户 | etk_… | Authorization + X-Tenant-Id | acme | 本租户内的资源管理、命令、事件、审计 | | 资源 | erk_… | Authorization + X-Tenant-Id + X-Resource-Id | acme/sensor-01 | 只有自己这一台 |
两种凭据的区分方式在各协议上是一致的:HTTP 看有没有 X-Resource-Id,MQTT 看 username 里有没有斜杠。所以资源请求必须始终带上 X-Resource-Id,MQTT username 必须是 tenantId/resourceId。
2.3 批量发证
一台一台调 /v1/resources 在出货时不现实。发证批次让资源自己换到密钥,控制台「凭据」页或 POST /v1/bootstrap 创建。四种模式对应不同的场景,共同点是结束状态一样:资源手上只有属于它自己的密钥,吊销一台就是一台。
| 模式 | 资源里放什么 | 约束 | 什么时候用 | | --- | --- | --- | --- | | open | 引导密钥 ebk_… | 无 | 试用。任意 ID、任意次数、可重复注册 | | limited | 引导密钥 ebk_… | 次数上限 + 有效期,每个 ID 一次 | 一般情况 | | roster | 引导密钥 ebk_… | 只认预登记清单,每个 ID 一次 | 出货,ID 提前知道 | | seed | 什么共享的都不放 | 出厂烧入 HMAC(种子, "租户/资源ID") | 安全第一,能控制产线 |
两道闸,因为它们的失效方式不同。 enabled 是个开关,状态看得见,但它只和「有没有人记得去关」一样可靠;expiresAt 会自己关,但十月设的日期到三月你不会记得。两个都有。有效期是绝对时间戳(expiresAt,Unix 毫秒)或 expiresInHours;设一个过去的时间点就等于立刻关闭,所以验证它不需要等、也不需要改系统时钟。
# 一般情况:100 台,72 小时窗口
curl -sS -X POST https://resources.example.com/v1/bootstrap \
-H "Authorization: Bearer $TENANT_API_KEY" -H 'X-Tenant-Id: acme' \
-H 'Content-Type: application/json' \
-d '{"bootstrapId":"line-a","mode":"limited","maxClaims":100,"expiresInHours":72}'
# → {"batch":{...},"apiKey":"ebk_..."} 只显示这一次
# 资源用它换自己的密钥
curl -sS -X POST https://resources.example.com/v1/enroll \
-H "Authorization: Bearer $BOOTSTRAP_KEY" -H 'Content-Type: application/json' \
-d '{"bootstrapId":"line-a","resourceId":"boiler-0001"}'
# → {"resourceId":"boiler-0001","apiKey":"erk_..."}几条必须知道的:
/v1/enroll在非 TLS 连接上直接 403。 换回来的是资源的长期凭据,明文传输的话窃听者拿到的东西比引导密钥还值钱——整套设计的收益会归零。本机测试可以ENROLL_REQUIRE_TLS=false,生产别关。- 一个资源 ID 只能领一次(
open除外)。第二次领返回409 already_claimed,而不是发一把新的顶掉在线的那台。这一条是整个机制成立的前提:直接注册接口是会覆盖已有凭据的,如果发证复用它,任何拿到引导密钥的人都能把在场资源踢下线。 - 被抢注是有声的。 真正那个资源装不上就是信号。清单里能看到
claimedAt;确认是误领或资源重置后,运维在控制台释放该 ID(DELETE /v1/bootstrap/{id}/claims?resourceId=…)才能再领。这是刻意做成运维动作——「让我再注册一次」正是攻击者会说的话。 - 关闭批次只挡新资源,已经发出去的密钥继续有效。要停某一台就删那台资源。
一机一密(seed)
资源那边没有任何共享物,也没有注册这一步——少掉的整个攻击面比多出来的强度更值钱:没有 enroll 端点可打、没有抢注竞争。
curl -sS -X POST .../v1/bootstrap -d '{"bootstrapId":"line-a","mode":"seed","prefix":"boiler"}'
# → {"batch":{...},"seed":"..."} 种子只给这一次,交给产线每个的密钥 = "erk_" + base64url(HMAC-SHA256(种子, "租户ID/资源ID"))。产线可以拿种子离线算,也可以调 POST /v1/bootstrap/line-a/derive 让平台算(同样要求 TLS)。资源直接用 acme/boiler-0001 + 这把密钥连接,首次连上时平台自动补一条资源记录。
prefix 限定这个种子能为哪些 ID 说话,别让一条产线的种子替另一条的资源背书。
两个代价得认下来:
- 种子是这套系统里唯一必须可还原存储的密钥,其他凭据都只存哈希。它用
CREDENTIAL_PEPPER派生的密钥做 AES-256-GCM 封装后入库,所以数据库被拷走不等于整批资源被拷走;但主机被完全攻破就没辙,而且丢了 pepper 就丢了种子。 - 换种子需要碰资源。 要轮换的话,正确做法是平台为每个资源重算它自己那一把新密钥、通过它已认证的通道下发——永远不要把种子发给资源。这部分还没实现。
3. HTTP/1.1 与 HTTP/2
端点:POST /v1/telemetry
curl --http2 -sS -X POST https://resources.example.com/v1/telemetry \
-H 'Authorization: Bearer erk_...' \
-H 'X-Tenant-Id: acme' \
-H 'X-Resource-Id: sensor-01' \
-H 'Content-Type: application/json' \
-d '{
"eventId":"018e8b28-3f3c-7ab9-a748-acde48001122",
"timestamp":1800000000000,
"data":{"online":true,"temperature":21.7,"humidity":45.2}
}'成功返回 HTTP 202。该响应表示可靠事件已经落盘,不只表示进入内存队列。管理员也可调用此端点,但请求体必须给出 tenantId 和 resourceId。
读取和续传:
curl -sS 'https://resources.example.com/v1/resources/sensor-01/events?tenantId=acme&afterSequence=128&limit=100' \
-H "Authorization: Bearer $ADMIN_API_KEY"如果游标已经落在保留窗口之外,结果会含 snapshotRequired: true。此时先读取 GET /v1/resources/{resourceId} 的最新快照,再从响应的 latestSequence 继续。
4. MQTT 与 MQTT over WebSocket
网关嵌入的 broker 支持 MQTT 3.1/3.1.1。
- TCP/TLS:
mqtts://resources.example.com:8883 - WebSocket:
wss://resources.example.com/mqtt,子协议mqtt
4.1 空间内就是普通 MQTT
一条连接落在一个 {域}/{空间} 里,由 username 决定;进去之后就是普通 broker:
- 主题完全是你的,网关不加任何层级,不写租户、不写空间
- 载荷随意,不要求 JSON
- 自收自发照常,订阅
#又发布tele/s,你会收到自己的消息
# 资源:主题自己定
mosquitto_pub -u 'acme/sensor-01' -P "$RESOURCE_API_KEY" -t 'tele/state' -m '{"t":21.7}'
mosquitto_pub -u 'acme/sensor-01' -P "$RESOURCE_API_KEY" -t 'resource/xx' -m '23.5'
mosquitto_pub -u 'acme/sensor-01' -P "$RESOURCE_API_KEY" -t 'tele/resource/xxx' -m 'ON'
# 后端:用你原有的 filter
mosquitto_sub -u 'acme' -P "$TENANT_API_KEY" -t 'tele/#' -v每个租户用自己的一套命名,互不影响;把资源标识放在 payload 里也完全可以——发送者身份由凭据决定,不靠主题携带,它在 $earth/events/… 的信封里,伪造不了。
载荷取值规则:
| 发出的载荷 | 记录下的 data | | --- | --- | | {"eventId":"…","timestamp":…,"data":{…}} | 幂等信封;重发同一个 ID 不会重复落盘 | | {"online":true,"t":21.7} | 原样,服务端生成 eventId | | 23.5 / ON / 任意纯文本 | {"value":23.5} / {"value":"ON"} |
只有第一种有去重。
4.2 username 决定域和空间
acme 租户密钥,本域,默认空间
acme/sensor-01 资源密钥,本域,默认空间
acme@ops 本域的空间 ops
acme@:ops 同上,域留空的写法
acme@ops,mkt 本域的两个空间;**第一个是主空间**
acme@* 本域里这个凭据能用的所有空间
acme@alliance: 组 alliance 的默认空间
acme@alliance:market 组 alliance 的 market 空间
acme@alliance:-,market 组 alliance 的两个空间
acme@alliance:* 组 alliance 里这个凭据能用的所有空间
acme/sensor-01@ops 资源密钥,本域的空间 ops(资源进不了组)冒号是「组」的唯一标志,连组的默认空间也要带冒号(acme@alliance:)。不带冒号的名字永远指你自己域里的空间,跟数据库里有什么无关。
一条连接持有多个空间
逗号分隔可以一次持有几个空间, 是「这个凭据能用的全部」。列表里的第一个是主空间, 时主空间是默认空间 -。
只持有一个空间时行为和以前逐字节一致,已经部署的客户端一个字都不用改。持有多个时:
tele/# 主空间。裸主题永远只是主空间,绝不因为多持有而多收
$space/ops/tele/# 指定空间
$space/+/tele/# 每个持有的空间里的这个主题
$space/# 每个持有的空间里的全部
$space/mkt/tele/s 发布,明确落在 mkt 一个空间投递时主空间保持裸名,其他空间带 $space/{空间}/ 前缀。这是必须的:两个空间里客户的主题完全一样,跨空间的订阅者否则分不清消息来自哪边。用 $ 前缀而不是普通层级,是因为普通第一层会和「客户自己就有 ops/ 开头的主题」撞名,而且 # 必须继续只表示主空间。
* 在 CONNECT 时就展开成具体列表,之后只做集合成员判断——许可集合永远来自凭据,不从客户端写的过滤器里解析。$space/ops/… 这种具名的直接编译成字面量;只有带 +/# 的才逐条投递时与集合求交。
对资源来说 * 是它注册时被授予的 spaces,不是整个域的——否则通配符就成了绕过逐个授权的方法。
集合里任何一个空间不允许,整条连接就被拒(CONNACK 5),不会「凭其中一个混进来」。
$space/ 下不能用通配符发布:MQTT 本身禁止 PUBLISH 主题带通配符,这里也一样。「发布到所有空间」没有拼写方式,这是故意的——扇出的发布收不回来,而且会让一个空间里出现没人在那里发过的消息。
被拒时说得清楚:
订阅 Subscription denied: this connection holds [ops], not "secret"
— name the spaces you need in the username, e.g. "user@ops,mkt" or "user@*"
发布 This connection may not publish into space "secret". It holds [ops]
— name the spaces you need in the username, e.g. "user@ops,mkt".这条规则看起来啰嗦,但换来的是一个名字在凭据的整个生命周期里只指一个地方。如果 acme@alliance 按「查一下是空间还是组」来解析,那么你哪天在自己域里建了个叫 alliance 的空间,所有在用 acme@alliance 的客户端就会被静默改道到另一个地方,两端都看不出来。
漏了冒号不会猜,会拒绝连接,并在日志里直接告诉你少了什么:
"g" is a group, not a space in "acme" — a group always takes a colon:
"@g:" for its default space, "@g:{space}" for one inside it把空间放在 username 而不是主题里,是为了让主题完全归你——已有客户端只改凭据,一个主题都不用动。一条连接对应一个空间;要同时用多个空间就开多条连接。
- 空间在控制台或
POST /v1/spaces创建,同租户内彼此隔离。 - 组是跨租户的域,由管理员创建并指定成员租户;成员在组的空间里能看到彼此的消息。资源进不了组——组是租户之间的约定。
- 组有自己独立的空间,和你自己域里的同名空间不是一回事。
acme@g:ops要的是g域里的ops;在acme域里建一个ops不管用。成员可以自己在组里建空间(控制台「空间」页的「我所在的组」里,或POST /v1/spaces带domain为组名);删除组内空间要管理员,因为它同时属于别的成员。 - 资源默认只能进默认空间,要进别的空间需在注册时用
spaces授予。
> 这一节讲规则。跑得出来的例子全部集中在主题、scope 与空间(例子版)——每种订阅写法实际收到什么、物理主题长什么样、被拒时会说什么,都在那里。
4.3 保留主题
$ 开头的主题归网关,客户端的 # 和 + 永远匹配不到它们——前提是 $ 在第一层。MQTT 3.1.1 §4.7.2 的原文是「以通配符开头的过滤器不匹配以 $ 开头的主题名」,所以我们把网关流量的物理名钉成 $earth/{域}/{空间}/…($ 在最前),而不是 {域}/{空间}/$earth/…——后者这条保护会无声消失。
$earth/message 是个例外,它是客户流量。 没有自带主题的命令投递到这里,物理形式是 {scope}/{域}/{空间}/$earth/message——$ 不在第一层,所以 # 收得到。也因此它可以被直接订阅:
mosquitto_sub -t '$earth/message' # 只收命令
mosquitto_sub -t '#' # 命令 + 空间里其他所有流量只想收命令的固件用前者,不必为了拿命令把整个空间的流量都拉过来。$earth/ 下其他名字仍然一律拒绝(SUBACK 0x80)并在 $earth/errors 上说明原因。
> 命令行里必须用单引号。 bash 会把 $earth 当未定义变量展开成空字符串,-t $earth/events/# 实际发出的是 -t /events/#,然后你会得到一片安静。写 -t '$earth/events/#'。网关会拒绝这种被吃掉 $ 的主题并说明原因。
| 主题 | 方向 | 用途 | | --- | --- | --- | | $earth/events/# | 网关 → 客户端 | 完整事件信封:eventId、sequence、权威 publisherId | | $earth/deliver/{publisherId}/{topic} | 客户端 → 单个客户端 | 记入可靠日志的定向投递,收件方在 {topic} 上收到 | | $earth/resume | 客户端 → 网关 | 补拉离线期间收到的消息 | | $earth/errors/# | 网关 → 客户端 | 被拒绝的发布/订阅及原因 |
4.4 定向投递与补拉
普通发布是空间内广播,不针对某个收件人,因此没有可补拉的收件通道。要让一条消息可审计、可重放,就得点名:
mosquitto_pub -u 'acme' -P "$TENANT_API_KEY" -q 1 \
-t '$earth/deliver/resource:sensor-01/cmd/relay1' -m '{"on":true}'资源在它自己订阅的 cmd/relay1 上收到,不需要额外分支。同一空间里的其他客户端收不到。
离线期间的消息不会自动补发(MQTT 会话是单实例内存,重启和换实例都会丢)。客户端连上后主动要回来:
mosquitto_pub -u 'acme/sensor-01' -P "$RESOURCE_API_KEY" -q 1 \
-t '$earth/resume' -m '{"afterSequence":128,"limit":100}'服务端把该序号之后、尚未过期、且不是自己发的消息按顺序重投到原主题,每条带 replay: true 和 eventId、sequence。按 eventId 去重即可——QoS 1 本来就会重投。
不要用 MQTT 会话状态做业务恢复,它跨不了实例,也扛不住重启。
4.5 被拒绝的发布
MQTT 在 QoS 0 下对发布没有任何回执,被拒的消息默认静默消失。网关会把原因发到 $earth/errors:
mosquitto_sub -u 'acme' -P "$TENANT_API_KEY" -t '$earth/errors/#' -v{"topic":"$earth/nope","error":"Unknown reserved topic …","at":1800000000000}用法错误(主题形状、载荷、未知保留名)报错但保持连接;只有越权(进不属于自己的空间、未授权的域)才断开。订阅被拒会直接返回 MQTT 订阅失败码,mosquitto_sub 当场打印 All subscription requests were denied.。
4.6 连接被拒
MQTT 3.1.1 的 CONNACK 带不了文字,只有几个返回码,所以连接失败的具体原因只在服务端日志里:
Connection Refused: bad user name or password → 凭据不对
Connection Refused: not authorised → 凭据没问题,但进不了这个域/空间后者有四种可能,日志会说明是哪一种:
reason: 'space "ops" does not exist in your domain "painterner" — create it with POST /v1/spaces'
reason: 'group "g" does not exist — create it with POST /v1/groups'
reason: '"painterner" is not a member of group "alliance" (members: other)'
reason: 'resource "s1" is granted [-], not "ops" — widen it by re-registering the resource with a `spaces` list'用 scripts/cluster-demo.sh logs a(或直接看网关日志)能看到。常见的一条是空间或组还没创建——acme@ops 里的 ops 必须先在控制台的「空间」页或 POST /v1/spaces 建出来,acme@g:ops 里的 g 还得是你所属的组。另一条常见的是进组漏了冒号——acme@g 指的是本域的空间 g,要组得写 acme@g:。
5. 通用 WebSocket
地址:wss://resources.example.com/v1/ws。Node、网关或 Agent 客户端在升级请求中发送:
Authorization: Bearer <apiKey>
X-Tenant-Id: acme
X-Resource-Id: sensor-01资源上报文本帧:
{"kind":"telemetry","eventId":"evt-000043","timestamp":1800000000000,"data":{"temperature":21.8}}持久化成功后返回:
{"kind":"telemetry-ack","version":1,"eventId":"evt-000043","sequence":129}订阅可靠事件使用 NetTrans core 协议:
{"kind":"subscribe","version":1,"channel":"tenant/acme/resource/sensor-01","afterSequence":128}处理每条 event 成功后发送 ACK:
{"kind":"ack","version":1,"channel":"tenant/acme/resource/sensor-01","sequence":129}服务端保证按序连续投递:订阅者收到的下一条一定是当前游标的下一个已落盘序号,不会跳号,也不会因为资源重试同一个 eventId 而收到重复事件。因此 sequence 可以当作「到此为止都已收到」的书签使用。
实现上,实时通道只被当作「有新数据」的信号:拿到的序号如果不是游标的下一个(并发写入可能让 12 先于 11 落盘),服务端不会直接下发,而是回到持久化日志按游标读取连续段——空洞在 RELIABLE_SEQUENCE_GAP_GRACE_MS 内会被等待,超时判定为废弃后跳过,通道不会永久卡住。
服务端拒绝超过已发送最高序号的 ACK,并把合法游标持久化。不要在收到事件、但业务处理尚未成功时提前 ACK。连接中断后用最后成功的 sequence 重新订阅。慢消费者超过 MAX_WS_BUFFERED_BYTES 时会以 1013 关闭,应退避重连。
浏览器原生 WebSocket 不能设置 Authorization 请求头。网关为此提供一次性票据:先用长期密钥调用 POST /v1/ws-tickets,再以 wss://resources.example.com/v1/ws?ticket=ewt_… 连接。票据默认 60 秒过期并在首次使用后销毁,长期密钥不进 URL。控制台和地球视图走的就是这条路径。
租户和管理员还可以订阅整个租户的实时流,用于仪表盘:
{"kind":"watch","tenantId":"acme"}之后每条持久化事件以 {"kind":"tenant-event","tenantId":"acme","event":{…}} 推送。这条通道没有游标、不重放、不需要 ACK——它服务于屏幕。要求不丢事件的消费端仍应订阅具体资源通道并按 sequence 确认。
6. HTTP/3 WebTransport
WebTransport 地址默认为 https://resources.example.com:4433/v1/webtransport,需显式设置 WEBTRANSPORT_ENABLED=true。连接后 10 秒内先发送控制请求:
const auth = await session.request('auth', {
username: 'acme/sensor-01',
apiKey: process.env.RESOURCE_API_KEY,
});
const ack = await session.request('resource.telemetry', {
eventId: crypto.randomUUID(),
timestamp: Date.now(),
data: { temperature: 21.9 },
});推荐客户端安装 @ptner/webtrans-lib-client 并使用 WebTransportResourceSession。控制信封运行在可靠双向流上,每帧为 4 字节大端长度加 JSON;命令以 resource.command notification 下发。events.read 请求可按 afterSequence 补拉。WebTransport 必须使用 HTTPS 和兼容证书。
WebTransport 是本服务的 HTTP/3 原生资源会话。普通 REST 由 Node 的 HTTP/2 端口提供;若业务要求 REST 本身通过 HTTP/3,可在前置 CDN/反向代理启用 HTTP/3,并回源到本服务的 HTTP/2。
7. 命令下发
管理员、租户或 MCP Agent 调用(资源密钥不能给其他资源下发命令):
curl -sS -X POST 'https://resources.example.com/v1/resources/sensor-01/commands?tenantId=acme' \
-H "Authorization: Bearer $TENANT_API_KEY" -H 'X-Tenant-Id: acme' \
-H 'Content-Type: application/json' \
-d '{"eventId":"cmd-42","expiresAt":1800000060000,"command":{"action":"sample","intervalSeconds":5}}'命令先写入可靠日志,再并行扇出到 MQTT 和已认证 WebTransport 会话。因此资源离线时 API 仍返回 202,资源可以在重连后通过事件回放获取命令。delivery 字段仅描述实时扇出尝试,不改变持久化结果。
下发到资源的命令载荷带有可靠信封的身份字段:
{
"eventId": "cmd-42",
"sequence": 129,
"tenantId": "acme",
"resourceId": "sensor-01",
"command": { "action": "sample", "intervalSeconds": 5 },
"expiresAt": 1800000060000,
"issuedAt": 1800000000000
}MQTT QoS 1 允许重复投递,重连恢复会话也可能重发,因此资源必须按 eventId 去重:记住最近执行过的 eventId,重复收到时确认但不重复执行。资源执行前还必须检查 expiresAt,过期命令应记录为忽略而不是执行。
8. 数据流识别(AI 判决器)
路由层不认识资源,这是刻意的:你想发什么主题就发什么主题。代价是平台本身也不知道哪条流是温度计、哪条是控制通道、哪条是谁在酒店大堂用 curl 试出来的。判决器补这个缺口——收到数据后送模型看一眼,打上标签。
三条它必须守住的规则:
- 不在投递路径上。 判决跑在持久化事件流之后。模型慢了、挂了、欠费了,消息照常投递。
- 不承担任何权限。 判决只是标签,不参与鉴权,打成
junk也不会丢数据。这一点是硬约束:被判决的 payload 是持有凭据的人写的,里面完全可能夹着冲模型说的话。 - 按「流」判决,不按「消息」。 判决单位是
{域, 空间, 发布者, 主题}加上 payload 的结构指纹。一个每秒发 10 条的传感器只是一个问题,不是每秒 10 个问题。
8.1 标签
| 标签 | 含义 | |---|---| | resource | 实体资源的遥测或状态(传感器、车辆、电表) | | resource | 普通业务记录,不是资源读数 | | control | 命令、设定值、执行动作 | | test | 探针、样例、试验流量 | | junk | 空的、坏的、对系统没有意义的 | | unknown | 没判出来(未启用、超限、判决失败) |
同时会尽量给出 resourceId 和一个 group 分组 slug。模型的输出一律当不可信文本校验:标签不在上表就算判决失败,resourceId 不像标识符就丢掉,绝不硬凑。
8.2 花多少钱是可控的
| 变量 | 默认 | 作用 | |---|---|---| | AI_JUDGE_API_KEY | 空 | 不配就整个关掉,标签仍可手动设置 | | AI_JUDGE_MODEL | claude-haiku-4-5-20251001 | 分类是最便宜的活,用小模型 | | AI_JUDGE_MAX_SAMPLE_BYTES | 2048 | 送去的样本上限;逐字段截断后再整体截断 | | AI_JUDGE_PER_TENANT_PER_MINUTE | 20 | 单租户频率 | | AI_JUDGE_PER_TENANT_PER_DAY | 500 | 单租户日限 | | AI_JUDGE_GLOBAL_PER_DAY | 5000 | 全站日限 | | AI_JUDGE_MIN_REJUDGE_MS | 3600000 | 同一结构多久内不再重判 | | AI_JUDGE_HEURISTICS | true | 明摆着的情况(payload 带经纬度、带 resourceId、主题里有 cmd)直接按规则判,不花钱 | | AI_JUDGE_CONCURRENCY | 2 | 并发调用数 | | AI_JUDGE_QUEUE_LIMIT | 500 | 待判队列上限,超了丢最旧的 |
配额是调用前预扣的,不是事后记账:一次发出去又失败的请求同样花了钱。超限不会静默跳过,会记成一条 source: "skipped" 的判决并写明原因,所以控制台里「这条为什么没标签」是有答案的。
启用它意味着把 payload 样本发给 AI_JUDGE_ENDPOINT 指定的模型服务商。这是一个关于客户数据去向的决定,所以默认关闭。
8.3 判错了怎么办
控制台「识别」页可以直接改,或者:
curl -X PUT "$BASE/v1/verdicts/$KEY?domain=acme" \
-H "authorization: Bearer $TENANT_API_KEY" -H 'content-type: application/json' \
-d '{"label":"resource","resourceId":"boiler-3","group":"temperature"}'人工设置的是最终答案,之后不会再被重判——一条你已经命名过的流,不该因为凌晨三点 payload 多了个字段就换标签。想让它重新判,DELETE 掉这条判决即可。
不设置也行:unknown 是个正常状态,系统没有任何地方依赖判决结果。
9. 错误与重试规则
400:信封或游标错误;修正后再发。401/403:凭据无效或越权;不要无限重试。404:资源尚未注册。413:超过请求或帧大小限制。429:若前置网关限流,按Retry-After退避。500/503、连接超时:指数退避,并使用同一eventId重发。
建议退避:1s、2s、4s、8s,加入 ±20% jitter,最高 60s。资源本地应保留尚未确认的有界队列,并监控队列深度与最老消息年龄。