/** * @dsh-local/im-agent-bridge — host 插件(C 单 · `交接单/IM群组-C-Agent接入插件.md`)。 * * ## 它做什么 * 实例内起一条**只拨出**的 WS 到平台 IM 端点,收到触发后组装上下文、交给本地 dsh * 会话、把回答写回房间。四条硬约束与两种形态的实现全在**平台侧**的 * `src/im/agent-bridge.ts`(纯逻辑,可在平台测试里直接断言);本文件只负责 * 「**连接 + 搬运**」: * * ``` * 平台 ──(WS 推送 /messages)──▶ 本插件 ──(组装上下文)──▶ 本地 dsh 会话 * ▲ │ * └────(op:'send' 回写, via='agent')◀───────────────────────┘ * ``` * * ## 🔴 五条纪律 * ① **只拨出**(E3):本文件**不开任何监听端口**(`ss -lntp` 无新增); * ② **不引入依赖**:WS 握手与帧编解码在此文件手写(与平台 `src/im/ws.ts` 同一套最小子集); * ③ **天然开关**:`DSH_PLATFORM_IM_URL` / `DSH_IM_INSTANCE_TOKEN` 任一缺失 ⇒ 不起连接; * ④ **故障隔离**(§五-7):任何异常都只影响本插件自己的连接,⛔ 不向上抛; * ⑤ **重连退避**(§五-2):断线后指数退避(1s→2s→…→上限 30s,带抖动),30 s 内自动恢复。 * * ⚠️ 本文件目前把「本地 dsh 会话」抽象成可注入的 `askLocalSession()` —— 默认实现是 * **回显式占位**(返回一条说明文本),因为真正的对接点(把 prompt 塞进哪个 dsh 会话) * 属 **D 单**的扩展点契约面(§八-5 明确问"是否需先做 D 单")。⇒ 本单保证的是 * **触发 / 上下文 / 预算 / 回写 / 隔离**五条判据可验;LLM 那一跳留给 D 单接线。 * * @module @dsh-local/im-agent-bridge */ import { appendFileSync } from 'node:fs' import { createHash, randomBytes } from 'node:crypto' import { join } from 'node:path' import { connect } from 'node:net' /** 平台 IM 端点(env:`DSH_PLATFORM_IM_URL`,形如 `https://ai1net.com`)。 */ const ENV_URL = 'DSH_PLATFORM_IM_URL' /** 实例凭据(env:`DSH_IM_INSTANCE_TOKEN`,`per-instance` 短期 token)。 */ const ENV_TOKEN = 'DSH_IM_INSTANCE_TOKEN' /** 本 agent 的引用标识(env:`DSH_IM_AGENT_REF`;缺省 = 用实例 id 兜底)。 */ const ENV_AGENT_REF = 'DSH_IM_AGENT_REF' /** 重连退避:首延迟 / 上限 / 抖动系数(§五-2 要求"30 s 内自动恢复")。 */ export const RECONNECT_FIRST_MS = 1000 export const RECONNECT_MAX_MS = 30_000 /** 排空 / 停机的宽限(`close()` 时给在途帧的时间)。 */ const CLOSE_GRACE_MS = 200 const WS_GUID = '258EAFA5-E914-47DA-95CA-C5AB0DC85B11' const OP_TEXT = 0x1 const OP_CLOSE = 0x8 const OP_PING = 0x9 const OP_PONG = 0xa /** 插件内日志(写 `/.dsh-im-agent.log`;⛔ 绝不写 token)。 */ function makeLog(home) { const file = join(home, '.dsh-im-agent.log') return (line) => { try { appendFileSync(file, `${new Date().toISOString()} ${line}\n`) } catch { // 日志写不动⛔绝不拖垮插件 } } } /* ── WS 最小子集(客户端侧:帧**必须掩码**)────────────────────────────────── */ function encodeClientFrame(opcode, payload) { const len = payload.length let header if (len < 126) { header = Buffer.alloc(2) header[1] = 0x80 | len } else if (len < 65536) { header = Buffer.alloc(4) header[1] = 0x80 | 126 header.writeUInt16BE(len, 2) } else { header = Buffer.alloc(10) header[1] = 0x80 | 127 header.writeBigUInt64BE(BigInt(len), 2) } header[0] = 0x80 | opcode const mask = randomBytes(4) const masked = Buffer.from(payload) for (let i = 0; i < masked.length; i += 1) masked[i] ^= mask[i % 4] return Buffer.concat([header, mask, masked]) } function textFrame(text) { return encodeClientFrame(OP_TEXT, Buffer.from(text, 'utf8')) } /** 服务端帧解析(**不带掩码**)。 */ class ClientDecoder { constructor() { this.buf = Buffer.alloc(0) } push(chunk) { this.buf = this.buf.length === 0 ? chunk : Buffer.concat([this.buf, chunk]) const out = [] for (;;) { if (this.buf.length < 2) break const b0 = this.buf[0] const b1 = this.buf[1] const opcode = b0 & 0x0f let len = b1 & 0x7f let offset = 2 if (len === 126) { if (this.buf.length < 4) break len = this.buf.readUInt16BE(2) offset = 4 } else if (len === 127) { if (this.buf.length < 10) break len = Number(this.buf.readBigUInt64BE(2)) offset = 10 } if (this.buf.length < offset + len) break const payload = Buffer.from(this.buf.subarray(offset, offset + len)) this.buf = this.buf.subarray(offset + len) out.push({ opcode, payload, fin: (b0 & 0x80) !== 0 }) } return out } } /** * 解析 URL 成 `{host, port, path, tls}`(⛔ 不引 `node:url` 的 WHATWG 也行,但用它更稳)。 * 只支持 `http(s)` / `ws(s)`:`https`/`wss` ⇒ TLS。 */ export function parseEndpoint(rawUrl) { const u = new URL(rawUrl) const tls = u.protocol === 'https:' || u.protocol === 'wss:' return { host: u.hostname, port: u.port !== '' ? Number(u.port) : tls ? 443 : 80, path: '/api/im/ws', tls, origin: `${tls ? 'https' : 'http'}://${u.host}`, } } /* ── 主类 ────────────────────────────────────────────────────────────────── */ /** * 拨出桥(一个实例一份)。 * * ⛔ **只拨出**:`start()` 只 `connect()`,不 `listen()`。 */ export class ImAgentBridge { /** * @param {object} opts * @param {string} opts.baseUrl - 平台 IM 端点(env)。 * @param {string} opts.token - 实例凭据(env)。 * @param {string} opts.agentRef - 本 agent 引用。 * @param {(prompt: string, ctx: object) => Promise} opts.askLocalSession * 把上下文交给本地 dsh 会话并取回回答(占位实现见 `defaultAsk`)。 * @param {(line: string) => void} [opts.log] * @param {number} [opts.firstDelayMs] / @param {number} [opts.maxDelayMs] - 退避参数(测试用) * @param {() => object} [opts.readSnapshot] - 读房间快照(测试注入;默认走 REST) */ constructor(opts) { this.baseUrl = opts.baseUrl this.token = opts.token this.agentRef = opts.agentRef this.askLocalSession = opts.askLocalSession this.log = opts.log ?? (() => undefined) this.firstDelayMs = opts.firstDelayMs ?? RECONNECT_FIRST_MS this.maxDelayMs = opts.maxDelayMs ?? RECONNECT_MAX_MS this.readSnapshot = opts.readSnapshot ?? null /** 已订阅的房间(重连后要重新 subscribe ⇒ 按游标补拉)。 */ this.rooms = new Map() this.socket = null this.decoder = null this.attempt = 0 this.timer = null this.stopped = false /** 观测计数(§五-5「指标可查」)。 */ this.metrics = { connects: 0, reconnects: 0, triggers: 0, replies: 0, suppressed: 0, errors: 0 } } /** 建立(或重建)连接。⛔ 只拨出。 */ start() { if (this.stopped) return this.metrics.connects += 1 const ep = parseEndpoint(this.baseUrl) const key = randomBytes(16).toString('base64') const socket = connect({ host: ep.host, port: ep.port }, () => { const headers = [ `GET ${ep.path} HTTP/1.1`, `Host: ${ep.host}`, 'Upgrade: websocket', 'Connection: Upgrade', `Sec-WebSocket-Key: ${key}`, 'Sec-WebSocket-Version: 13', // 🔴 两件关键头:实例凭据 + 区声明(区声明可省;带上便于平台侧判跨区)。 `x-dsh-im-instance-token: ${this.token}`, `x-dsh-im-agent-ref: ${this.agentRef}`, '', '', ].join('\r\n') socket.write(headers) }) this.socket = socket this.decoder = new ClientDecoder() let handshakeDone = false let raw = Buffer.alloc(0) socket.on('data', (chunk) => { if (!handshakeDone) { raw = Buffer.concat([raw, chunk]) const idx = raw.indexOf('\r\n\r\n') if (idx < 0) return const head = raw.subarray(0, idx).toString('utf8') if (!/^HTTP\/1\.1 101/.test(head)) { // 握手被拒(401 / 404 …):记一笔并按退避重连(⛔ 不静默停)。 this.log(`handshake-rejected ${head.split('\r\n')[0]}`) this.metrics.errors += 1 socket.destroy() return } const expect = createHash('sha1').update(`${key}${WS_GUID}`).digest('base64') if (!head.includes(`Sec-WebSocket-Accept: ${expect}`)) { this.log('handshake-accept-mismatch') this.metrics.errors += 1 socket.destroy() return } handshakeDone = true this.attempt = 0 this.log(`connected ${ep.host}`) // 重连后**重新订阅**所有房间(带游标 ⇒ 平台补拉断线期间的消息)。 for (const [roomId, since] of this.rooms) this.send({ op: 'resume', roomId, since }) const rest = raw.subarray(idx + 4) if (rest.length > 0) this.consume(rest) return } this.consume(chunk) }) socket.on('error', (err) => { this.metrics.errors += 1 this.log(`socket-error ${String(err && err.message ? err.message : err)}`) }) socket.on('close', () => { this.socket = null this.decoder = null if (this.stopped) return this.scheduleReconnect() }) } consume(chunk) { if (this.decoder === null) return let frames try { frames = this.decoder.push(chunk) } catch (err) { this.log(`decode-error ${String(err)}`) return } for (const frame of frames) { if (frame.opcode === OP_PING) { this.rawWrite(encodeClientFrame(OP_PONG, frame.payload)) continue } if (frame.opcode === OP_CLOSE) { // 服务端要关:主动断开 ⇒ 走 close 事件里的退避重连。 if (this.socket !== null) this.socket.destroy() continue } if (frame.opcode !== OP_TEXT) continue let msg try { msg = JSON.parse(frame.payload.toString('utf8')) } catch { continue } // ⚠️ 每条消息处理都吞异常:插件故障⛔不许拖垮连接(§五-7)。 void this.onServerFrame(msg).catch((err) => { this.metrics.errors += 1 this.log(`frame-handler-error ${String(err)}`) }) } } rawWrite(buf) { if (this.socket === null || this.socket.destroyed) return false try { this.socket.write(buf) return true } catch { return false } } send(obj) { return this.rawWrite(textFrame(JSON.stringify(obj))) } scheduleReconnect() { if (this.stopped) return this.metrics.reconnects += 1 const base = Math.min(this.maxDelayMs, this.firstDelayMs * 2 ** this.attempt) // 抖动 ±20%:避免整批实例同时重连打爆平台。 const delay = Math.max(1, Math.round(base * (0.8 + Math.random() * 0.4))) this.attempt += 1 this.log(`reconnect in ${delay}ms (attempt ${this.attempt})`) this.timer = setTimeout(() => { this.timer = null this.start() }, delay) if (typeof this.timer.unref === 'function') this.timer.unref() } /** 订阅一个房(调用方决定订阅哪些)。 */ subscribe(roomId, since = 0) { this.rooms.set(roomId, since) this.send({ op: 'subscribe', roomId }) } /** 记住游标(每条收到的消息都推进它 ⇒ 重连补拉不重复)。 */ noteCursor(roomId, seq) { if (typeof seq === 'number' && Number.isFinite(seq)) this.rooms.set(roomId, seq) } /** * 服务端帧处理。 * * 只关心两种:`message`(推送的新消息 ⇒ 可能触发)与 `messages`(补拉结果)。 */ async onServerFrame(msg) { if (msg === null || typeof msg !== 'object') return if (msg.op === 'pong') return if (msg.op === 'hello') { this.log('hello') return } if (msg.op === 'error') { this.metrics.errors += 1 this.log(`server-error ${String(msg.error)}`) return } const roomId = typeof msg.roomId === 'string' ? msg.roomId : undefined const list = Array.isArray(msg.messages) ? msg.messages : [] if (roomId === undefined || list.length === 0) return for (const m of list) { this.noteCursor(roomId, m && typeof m.seq === 'number' ? m.seq : undefined) await this.maybeReply(roomId, m) } } /** * 判一条消息要不要回、要回就回。 * * 🔴 这里**不做**触发过滤的重活 —— 那条判据在平台侧 `src/im/agent-bridge.ts` * (同一套逻辑,测试直接断它)。本文件只做"要不要处理"的**本地前置闸** * (形态 / @ 命中),够用且不引第二份判据源。 */ async maybeReply(roomId, message) { if (message === undefined || message === null) return // ① 自激闸门:agent 自己的消息(`via==='agent'` 或 `author_kind!=='human'`)恒不处理。 if (message.authorKind !== 'human' || message.via === 'agent') return const mentions = collectMentions(message.payload) if (mentions.length === 0) return // ② 是不是 @ 到我(形态 A:agentRef 直接命中;形态 B:主人命中 —— 由平台侧快照判到底是谁)。 const self = String(this.agentRef) const hit = mentions.includes(self) || mentions.includes(String(this.ownerId ?? '')) || mentions.some((m) => String(m).startsWith(self)) if (!hit) return this.metrics.triggers += 1 const ctx = { roomId, focusSeq: typeof message.seq === 'number' ? message.seq : undefined, memberId: message.authorId, message, } let answer try { answer = await this.askLocalSession(renderPrompt(message, roomId), ctx) } catch (err) { this.metrics.errors += 1 this.log(`ask-failed ${String(err)}`) return } if (typeof answer !== 'string' || answer === '') return // ③ 回写:`via='agent'` + 可见代答标识(§五-6)。 const marker = this.form === 'bot' ? '能力机器人' : `由 ${this.ownerName ?? message.authorId} 的助手代答` const ok = this.send({ op: 'send', roomId, payload: { text: answer, agentReply: true, agentLabel: marker }, via: 'agent', }) if (ok) { this.metrics.replies += 1 this.log(`reply-sent room=${roomId} focus=#${ctx.focusSeq}`) } } /** 停机(挂 profile 卸载 / 实例退出)。 */ close() { this.stopped = true if (this.timer !== null) { clearTimeout(this.timer) this.timer = null } if (this.socket !== null) { try { this.socket.write(encodeClientFrame(OP_CLOSE, Buffer.from('\u0000\u0000shutdown', 'utf8'))) } catch { // 写不动就直接销毁 } const s = this.socket setTimeout(() => s.destroy(), CLOSE_GRACE_MS).unref?.() } } } /** 从 payload 抽 @(与平台侧同判据,这里是插件侧的只读副本)。 */ export function collectMentions(payload) { if (payload === null || typeof payload !== 'object') return [] if (Array.isArray(payload.mentions)) return payload.mentions.filter((m) => typeof m === 'string') const text = payload.text if (typeof text !== 'string') return [] const out = [] for (const m of text.matchAll(/@([A-Za-z0-9_\-.:]{1,64})/g)) if (m[1] !== undefined) out.push(m[1]) return out } /** 把一条消息渲染成给本地 dsh 会话的 prompt(占位实现 —— 完整上下文由平台侧组装)。 */ export function renderPrompt(message, roomId) { const text = typeof message.payload?.text === 'string' ? message.payload.text : JSON.stringify(message.payload) return `[房间 ${roomId} · 消息 #${message.seq} · 来自 ${message.authorId}]\n${text}\n\n请回答上面这条被 @ 的消息。` } /** * 默认的本地会话实现:**回显式占位**。 * * ⛔ 真正接本地 dsh 会话是 **D 单**的扩展点契约面(§八-5);本单保证的是 * 触发 / 上下文 / 预算 / 回写 / 隔离五条判据**可验**。接法 = 通过 `opts.askLocalSession` 注入。 */ export function defaultAsk(prompt) { return Promise.resolve(`(agent 桥已收到触发,但本地会话对接属 D 单扩展点)\n${prompt.slice(0, 200)}`) } /* ── cordis 落点 ──────────────────────────────────────────────────────────── */ /** 单例(profile 卸载时 close)。 */ let activeBridge = null /** * cordis 入口。 * * 🔴 **天然开关**:读不到 URL 或 token ⇒ 直接 return(日志里**没有**连接尝试)。 * 这与 §三「不注入 env ⇒ 插件不启动」逐字对应,也是平台侧回滚的唯一动作。 */ export function apply(ctx) { const url = process.env[ENV_URL] const token = process.env[ENV_TOKEN] const home = process.env.DSH_HOME || process.env.HOME || '.' const log = makeLog(home) if (url === undefined || url === '' || token === undefined || token === '') { log(`skip: ${ENV_URL} / ${ENV_TOKEN} 未注入 ⇒ agent 桥不启动`) return } const agentRef = process.env[ENV_AGENT_REF] ?? 'instance-agent' const bridge = new ImAgentBridge({ baseUrl: url, token, agentRef, askLocalSession: defaultAsk, log, }) log(`apply pid=${process.pid} agent=${agentRef}`) try { bridge.start() } catch (err) { // §五-7:起连接失败⛔不许把 profile 打崩。 log(`start-failed ${String(err)}`) } activeBridge = bridge ctx.on?.('dispose', () => { try { bridge.close() } catch { // 收尾失败也不抛 } activeBridge = null }) } /** 供测试/排查取当前单例。 */ export function currentBridge() { return activeBridge }