/** * device lane **WS 汇聚端**(design/device-executor-lane-v2 §5.2 连接状态机 / §5.3 双层超时 / * §5.4 重连与幂等 / §5.5 背压 / §8-R7 发版排空)—— 本仓**首个** WS 服务端。 * * ── 三条它必须同时守住的不变量 ─────────────────────────────────────────────────────────────────── * ① **at-most-once**:一条指令至多在设备上执行一次。`dispatch_committed` 是**不可回退的单向门**, * 且在**首次 socket 写之前**原子落定 —— `ws.send()` 成功只证明字节进了本机缓冲。commit 之后的 * 一切不确定恒收敛到 `outcome_unknown`,**永不重发**(§5.3/§5.4)。 * ② **世代闸**:设备侧**所有**可变状态帧携 `gen`,判据 = **严格等于**本 socket 协商世代(非 ≥)。 * 旧世代的 result/chunk/heartbeat/goodbye 在进任何缓冲之前即被拒 + 审计(D-P4 四格)。 * ③ **结果帧四不变量**(§5.1):已鉴权且 deviceId 与 pending 行一致 / instructionId 在本进程 pending * 表且态 = `dispatch_committed` / `deviceSeq` 单调(重复 = 幂等吸收) / `streamSeq` 连续 * (乱序缓冲,缺片 = 指令失败**而不是**静默截断)。 * * ── A-2 → A-3 的移交义务(工作稿「着陆实况」三条,逐条兑现在哪儿)──────────────────────────────── * ① **下发边界必须重新断言世代**(codex F4):`admitSession` 给的 `connGeneration` 只是一个**快照**, * 它与真正写 socket 之间可以插进一次 `revokeDevice`。所以 {@link createDeviceWsHub} 的 dispatch 在 * 写 socket 之前做三重断言:socket 世代 === 判决世代、本 conn 仍是该设备的当前连接、以及一次 * **带世代谓词的库写**(`heartbeatConnection`)—— 影响 0 行 = 本连接已失权 ⇒ 停发 + `bye` + 断连。 * 🔴 **诚实残余**:这三条关不掉「断言通过之后、`ws.send` 之前」的那一瞬。窗口上界 = 租约续期拍距 * (`leaseRenewIntervalMs`,默认 = 心跳周期):吊销后最迟一拍,连接就会在续租里发现失权并断开。 * 窗内已发出的指令按 §8-R4 记账 —— 在途结果帧丢弃 + 审计,对应调用按 outcome-unknown 收尾。 * 真正的最终防线在**设备侧**(§5.4 ⑦ 单实例锁 + 严格等值世代判据),云侧关不了这一格。 * ② 事务内回读同连接(`readDeviceOn`)是**店内**规矩;本车不新增任何 SQL,只消费店契约。 * ③ 消费 A-2 预留的两个码(`device.busy` / `device.attached_elsewhere`)与两个审计词 * (`device_frame_unexpected` / `device_superseded`)—— 见下方各自的落点。 * * ── 为什么 upgrade 处理器必须显式处理「不认的路径」──────────────────────────────────────────────── * node 的默认行为是「没有 'upgrade' 监听器就销毁 socket」。挂上第一个监听器之后那条默认**失效**, * 于是不认的路径若直接 `return`,socket 会挂着直到超时。所以这里显式回 404 并销毁。 * ⚠️ 将来出现**第二个** upgrade 消费方时,这里要改成一张路由表 —— 两个 `server.on("upgrade")` * 监听器会**同时**收到每一次 upgrade,各自「不认就 404」会互相把对方的连接打死。 * * ── 安全轴姿势([ref])──────────────────────────────────────────────────────────────────────────── * 未知帧型 / 坏形帧 / 世代不符 / pending 外 instructionId = **响亮拒 + 审计行**,没有一处空 catch, * 没有一处「认不出就当默认值」。审计写失败本身也是响亮的(logger.error + 计数),不是静默吞。 */ import http from "node:http"; import { type DeviceOwner, type DeviceStore } from "./device-store.js"; import type { DeviceEnrollment, DeviceRejection } from "./device-enrollment.js"; import { type DeviceChunkKind, type DeviceInstructionArgsByKind, type DeviceInstructionKind, type DeviceInstructionOutcome } from "./device-ws-protocol.js"; /** 交给调用方的一片上行流(`data` 已从 base64 解回字节)。 */ export interface DeviceChunkDelivery { streamSeq: number; kind: DeviceChunkKind; data?: Buffer; exitCode?: number; } /** * 一次下发的入参。`kind` 与 `args` **按 kind 判别**(§5.1 task-R1-8),不是 `any`。 * * `payload` = 大载荷下行(writeFile 族):内容不进指令帧,走 `payload` begin→chunk→commit 子序; * 长度与 sha256 由**本层**按真字节重算写回 args —— 调用方报的摘要不作数(它不是那份字节的权威)。 */ export type DeviceDispatchInput = { [K in DeviceInstructionKind]: { rootSessionId: string; requester: DeviceOwner; /** * O4 §3 投递门**身份重断言**(codex 轮2 R2-[high],红先 R4b 实证):调用方 env **铸造时**钉住的 * 设备。准入在 dispatch 内解出的是**当前**绑定 —— 显式换绑与在飞 turn 竞态时两者可以不同,而 * adapter 的 connect 身份检查罩不住已连接 env 的后续调用(它们不再过 connect)。差异 ⇒ 恒拒 * `device.identity_mismatch`(cwd/能力面是从铸造设备读的,静默改道 = 在错误机器的工作区上执行)。 * **必填**:可选会让漏传静默([ref] `originTaskId` 同判)。 */ expectedDeviceId: string; kind: K; args: DeviceInstructionArgsByKind[K]; cwd?: string; shellEnv?: Record; timeoutMs: number; payload?: Uint8Array; payloadPartBytes?: number; onChunk?: (chunk: DeviceChunkDelivery) => void; /** * 铸造期快照的**否定面**(codex R3-F1):调用方在铸自己那一刻**没有**看到这几个指令 kind, * 所以它的对外面上没有对应的方法。若投递门这一刻的连接**声明了**其中任何一个 ⇒ 拒。 * * 🔴 为什么这条判据非放在这里不可:调用方侧那一次检查与真正的下发之间隔着 `admitSession` 的一次 * SQL 往返 —— 设备可以在这个窗口里断开、以**更强**的能力连回来,于是检查看到的是旧连接、指令却 * 打到新连接上(TOCTOU)。而投递门是本进程里**唯一**同时持有「最终选中的连接」与「不可回退的 * commit」的地方,原子复核只能在这儿做。语义与世代重断言同族(A-2 移交义务①),放在同一段。 * * 用途在 device lane 上是**保护型写原语**(`writeFileGuarded`/`writeFileExclusive`):它们的缺席被 * core 读成「后端没有这个原语」并降级成非原子写,所以「铸的时候没有、下发的时候有了」= 一次静默的 * 保护降级([ref])。本座本身与语义无关,只做集合判断。 */ refuseIfDeclared?: readonly DeviceInstructionKind[]; /** * 车A-4(A-3 残余⑤):调用方的取消轴。 * * 🔴 **为什么这个座必须在 hub 而不是调用方自己接**:`instructionId` 是在**投递门内部**铸的, * `dispatch` 直到结算才 resolve —— 调用方结构上拿不到那个 id,所以它自己的 abort 监听只能做到 * 「放弃等待」,而设备上那条命令还在跑(§5.4 at-most-once 的反面:一次不可见的执行)。座放在这里, * abort ⇒ 真发 `cancel` 帧,终局回执照常走 §5.3 的 10s 宽限。 * 投递门**之前** abort ⇒ `terminal_never_started`(设备结构上没见过它,安全)。 */ signal?: AbortSignal; }; }[DeviceInstructionKind]; /** * 下发结果的三态 —— 与 {@link DeviceInstructionState} 的两条路径 1:1。 * * • `settled`:过了投递门(`dispatch_committed`)并拿到终局。**含**失败结局(设备回的 error / * 云侧兜底的 `outcome_unknown`)—— 它们都是「设备可能已经执行过」的世界。 * • `never_started`:`terminal_never_started` —— 设备结构性不可能见过它(能力未声明)。可安全重试。 * • `rejected`:准入/背压/失权,**在投递门之前**。同样安全(§5.3 投递边界)。 */ export type DeviceDispatchResult = { kind: "settled"; instructionId: string; outcome: DeviceInstructionOutcome; } | { kind: "never_started"; outcome: DeviceInstructionOutcome; } | { kind: "rejected"; reject: DeviceRejection; }; export interface DeviceWsHubLogger { warn(msg: string, fields?: Record): void; error(msg: string, fields?: Record): void; info?(msg: string, fields?: Record): void; } export interface DeviceWsHubOptions { store: DeviceStore; /** 准入链(§4.3.2)。WS 汇聚端**不自建判据**,一律经它。 */ enrollment: DeviceEnrollment; /** 本副本 id(写进连接租约行,跨副本排障用)。 */ replicaId: string; logger?: DeviceWsHubLogger; metrics?: { inc(name: string, labels?: Record): void; }; /** 下发给设备的心跳周期(helloAck 协商);丢 `heartbeatMisses` 拍判 STALE。 */ heartbeatIntervalMs?: number; heartbeatMisses?: number; /** 连接租约时长(整秒,自 DB now 起算)。 */ leaseSeconds?: number; /** 服务端续租/失权侦测的拍距。缺省 = `heartbeatIntervalMs`(见文件头「诚实残余」)。 */ leaseRenewIntervalMs?: number; /** per-device 在途上限(§11-O3 clay 已裁 = 4)。 */ maxInflightPerDevice?: number; /** 等在途槽的上限;到期 = `device.busy`(§5.5,429 语义)。 */ dispatchTimeoutMs?: number; /** 云侧执行超时兜底的宽限(§5.3:`timeoutMs + grace` ⇒ outcome_unknown)。 */ execGraceMs?: number; /** cancel 终局回执宽限(§5.3)。 */ cancelGraceMs?: number; /** 断连后在途不判丢的窗(§5.4 同 epoch 重连);窗尽 ⇒ outcome_unknown。 */ reconnectGraceMs?: number; /** 上行流片缺口的等待上限;到期 = 指令失败(**不是**静默截断)。 */ chunkGapTimeoutMs?: number; /** §5.2 顶替旋钮:`stale_only`(默认)= 只顶替已 STALE 的旧连;`deny` = 重连也须等旧连宽限尽。 */ supersede?: "stale_only" | "deny"; /** codex R1-F5:per-instruction 乱序缓冲的条目帽(超出 = 指令失败 + 审计)。 */ chunkBufferMaxEntries?: number; /** codex R1-F5:per-instruction 乱序缓冲的**解码后**字节帽(§5.5 的 2MiB)。 */ chunkBufferMaxBytes?: number; /** codex R1-F3:本副本**未鉴权**连接的并发帽(超额在 upgrade 阶段就 503 拒)。 */ maxPreAuthConnections?: number; /** codex R1-F3:单来源地址的未鉴权连接并发帽。[ref]([ref]):字面 `off` = 显式停用该层—— * 帽键是裸 remoteAddress,XFF 反代 / 公司 NAT 部署形把全员折叠成单键,这层从公平层退化成 * 「全员共享一份预算」;停层后全局帽(`maxPreAuthConnections`)照常兜底。 */ maxPreAuthConnectionsPerRemote?: number | "off"; } export interface DeviceWsHub { /** 把 `GET /v1/device/ws` 的 upgrade 处理器挂到裸 `http.Server` 上。 */ attach(server: http.Server): void; dispatch(input: DeviceDispatchInput): Promise; /** 幂等:重复 cancel 同一指令 = no-op(§5.3)。 */ cancel(instructionId: string): Promise; /** §8-R7:停发新指令 + 拒新 upgrade,但**保持既有 WS**,在途结果帧照收。 * [ref]:装配层可带 `graceMs`(= shutdown 的 drainGraceMs 上界)——排空拒体借它给出「drain 窗最迟 * 何时收窗」的指路;getter 形每次活读(drainGraceMs 是热改键,一次性快照会在窗中热调后失真, * codex-S09-F2);不带 = 拒体不造数(测试宿主/嵌入形没有这个事实)。 */ beginDrain(opts?: { graceMs?: number | (() => number); }): void; isDraining(): boolean; connectedDeviceIds(): string[]; /** * S-470(车C C-G2「吊销后活连接 5s 内断」)—— **立刻**对这台设备的活连接重跑一次世代围栏。 * * 🔴 为什么是「重跑围栏」而不是一条新的断连路径:吊销已经在**店内同事务**作废租约并把 * `conn_generation` 推进一格(§4.6 步 2),所以本连接的下一次带世代谓词的写必然影响 0 行 —— * {@link fenceConnection} 于是发 `bye:revoked`、断连、把在途按 outcome-unknown 收尾。管理面要的 * 「立刻」只是**把那一次围栏提前到现在**,不是第二种断法(两种断法 = 两份 bye 语义会漂)。 * 定时器驱动的续租(`leaseRenewIntervalMs`)仍是**跨副本**的兜底:连接不在本副本时本调用是 no-op, * 那台设备最迟在一个续租拍距内失权。 */ fenceDevice(deviceId: string): Promise; /** * 车A-4:该设备**当前连接**在 hello 里声明的指令面(与本 server 的闭集取过交集之后的那一份)。 * 没有活连接 ⇒ `undefined`(**不是**空集:「不知道」与「明确什么都不支持」是两件事)。 * * 🔴 用途只有一个,而且是 core 契约逼出来的:`writeFileExclusive`/`writeFileGuarded` 是 * **presence-typed** 可选面 —— core 的降级律逐字要求「后端没有该原语就把方法留 undefined,禁两步 * 仿真」,所以「在不在场」必须在**铸 env 的那一刻**答得出来,不能等到第一次调用再回 not_supported * (那时调用方已经按「有原子面」走了)。 * 诚实残余:设备**在铸 env 之后**才连上时,这一问答的是 undefined ⇒ 该面缺席 ⇒ 写门如实降级成 * 「advisory adjudication」而不是宣称关闭了 TOCTOU —— 方向是安全的那一侧。 */ declaredSupports(deviceId: string): ReadonlySet | undefined; close(): Promise; } export declare function createDeviceWsHub(options: DeviceWsHubOptions): DeviceWsHub; //# sourceMappingURL=device-ws-hub.d.ts.map