| 1 | import test from "node:test"; |
| 2 | import { createServer } from "node:http"; |
| 3 | import { once } from "node:events"; |
| 4 | import assert from "node:assert/strict"; |
| 5 | import { mkdtemp, readFile, rm } from "node:fs/promises"; |
| 6 | import { tmpdir } from "node:os"; |
| 7 | import path from "node:path"; |
| 8 | import url from "node:url"; |
| 9 | import { |
| 10 | ApprovalOwnershipError, ThreadStore, decideApproval, settleStoredApproval, |
| 11 | createRuntimeClient, activeTurnBlock, readSse, readJsonSafe, compactRuntimeError |
| 12 | } from "../src/lib.mjs"; |
| 13 | import { incomingIdentity as wecomIdentity } from "../../wecom-bridge/src/lib.mjs"; |
| 14 | |
| 15 | const root = path.resolve(path.dirname(url.fileURLToPath(import.meta.url)), "..", ".."); |
| 16 | const chatId = "chat-a", threadId = "thread-a", turnId = "turn-a"; |
| 17 | const owner = { chatId, threadId, turnId, actorId: "telegram-user:alice" }; |
| 18 | const identities = { |
| 19 | "telegram-bridge": (userId) => ({ chatId, chatType: "group", userId, isBot: false }), |
| 20 | "feishu-bridge": (openId) => ({ chatId, chatType: "group", openId }), |
| 21 | "wecom-bridge": (userId) => ({ body: { chatid: chatId, chattype: "group", from: { userid: userId } } }), |
| 22 | "weixin-bridge": () => undefined |
| 23 | }; |
| 24 | const actors = { |
| 25 | "telegram-bridge": "telegram-user:alice", "feishu-bridge": "feishu-openId:alice", |
| 26 | "wecom-bridge": "wecom-user:alice", "weixin-bridge": "weixin-account:bot-a:user:chat-a" |
| 27 | }; |
| 28 | async function withStore(fn, options = {}) { |
| 29 | const dir = await mkdtemp(path.join(tmpdir(), "bridge-approvals-")); |
| 30 | try { |
| 31 | const store = await ThreadStore.open(path.join(dir, "state.json"), { actions: true, ...options }); |
| 32 | await store.setChat(chatId, { threadId }); |
| 33 | await fn(store); |
| 34 | } finally { await rm(dir, { recursive: true, force: true }); } |
| 35 | } |
| 36 | function extract(source, name) { |
| 37 | const marker = new RegExp(`^(?:async )?function ${name}\\(`, "m"); |
| 38 | const start = source.search(marker); |
| 39 | assert.notEqual(start, -1, name); |
| 40 | const rest = source.slice(start); |
| 41 | const next = rest.search(/\n(?:async )?function \w+\(/); |
| 42 | return next === -1 ? rest : rest.slice(0, next); |
| 43 | } |
| 44 | async function shipping(bridge, names, values) { |
| 45 | const source = await readFile(path.join(root, bridge, "src/index.mjs"), "utf8"); |
| 46 | const functions = bridge === "weixin-bridge" && names.includes("turnActor") ? ["accountIdentity", ...names] : names; |
| 47 | return new Function(...Object.keys(values), functions.map((name) => extract(source, name)).join("\n") + |
| 48 | `\nreturn { ${names.join(", ")} };`)(...Object.values(values)); |
| 49 | } |
| 50 | async function withRuntime(fn) { |
| 51 | const state = { pending: [{ id: "approval-a", turn_id: turnId }], turns: [], posts: [], |
| 52 | startAttempts: 0, accepted: [], failStart: false, acceptDecision: true, events: [] }; |
| 53 | const server = createServer(async (request, response) => { |
| 54 | assert.equal(request.headers.authorization, "Bearer bridge-runtime-fixture-token"); |
| 55 | let text = ""; |
| 56 | for await (const chunk of request) text += chunk; |
| 57 | response.setHeader("connection", "close"); |
| 58 | if (request.url.includes("/events?")) { |
| 59 | response.setHeader("content-type", "text/event-stream"); |
| 60 | response.end(state.events.map((event) => `data: ${JSON.stringify(event)}\n\n`).join("")); |
| 61 | return; |
| 62 | } |
| 63 | response.setHeader("content-type", "application/json"); |
| 64 | if (request.method === "POST" && request.url.endsWith("/turns")) { |
| 65 | state.startAttempts++; |
| 66 | if (state.failStart) { |
| 67 | response.statusCode = 503; |
| 68 | response.end(JSON.stringify({ error: { message: "start rejected" } })); |
| 69 | } else { |
| 70 | const id = `accepted-${state.startAttempts}`; |
| 71 | state.accepted.push(id); |
| 72 | response.end(JSON.stringify({ turn: { id } })); |
| 73 | } |
| 74 | } else if (request.method === "POST") { |
| 75 | state.posts.push({ route: request.url, body: JSON.parse(text) }); |
| 76 | response.statusCode = state.acceptDecision ? 200 : 503; |
| 77 | response.end(JSON.stringify(state.acceptDecision ? { ok: true } : { error: { message: "fixture unavailable" } })); |
| 78 | } else { |
| 79 | response.end(JSON.stringify({ id: threadId, pending_approvals: state.pending, turns: state.turns, latest_seq: 0 })); |
| 80 | } |
| 81 | }); |
| 82 | server.listen(0, "127.0.0.1"); |
| 83 | await once(server, "listening"); |
| 84 | const runtimeUrl = `http://127.0.0.1:${server.address().port}`; |
| 85 | const runtime = createRuntimeClient({ runtimeUrl, runtimeToken: "bridge-runtime-fixture-token" }); |
| 86 | try { await fn({ ...runtime, runtimeUrl, state }); } |
| 87 | finally { await new Promise((resolve) => server.close(resolve)); } |
| 88 | } |
| 89 | function environment(store, runtime, extra = {}) { |
| 90 | return { |
| 91 | threadStore: store, runtimeJson: runtime.runtimeJson, decideRuntimeApproval: decideApproval, |
| 92 | settleStoredApproval, ApprovalOwnershipError, incomingIdentity: wecomIdentity, |
| 93 | sendText: async () => {}, replyText: async () => {}, sendTurnText: async () => {}, |
| 94 | parseApprovalDecisionArgs: () => { throw new Error("fixture uses actual parsed approval id"); }, |
| 95 | ensureThread: async () => store.getChat(chatId), activeTurnBlock, |
| 96 | config: { runtimeUrl: runtime.runtimeUrl, turnTimeoutMs: 1000, maxReplyChars: 10000 }, |
| 97 | controlKeyboard: () => ({}), activeTurnKeyboard: () => ({}), helpText: () => "", |
| 98 | clearActiveTurn: () => store.patchChat(chatId, { activeTurnId: null }), stopping: false, |
| 99 | streamTurnEvents: async (chat, thread, turn) => { |
| 100 | const reopened = await ThreadStore.open(store.filePath, { actions: true }); |
| 101 | assert.ok(reopened.turnOrigin(chat, thread, turn), "origin must be durable before delivery"); |
| 102 | }, ...extra |
| 103 | }; |
| 104 | } |
| 105 | async function invokeStart(bridge, fn, user = "alice") { |
| 106 | if (bridge === "telegram-bridge") return fn(chatId, "prompt", { actorId: `telegram-user:${user}` }); |
| 107 | return fn(chatId, "prompt", identities[bridge](user)); |
| 108 | } |
| 109 | |
| 110 | test("approval tokens retain 128-bit entropy and require their accepted human and turn", async () => { |
| 111 | await withStore(async (store) => { |
| 112 | await store.recordTurnOrigin(chatId, threadId, turnId, owner.actorId); |
| 113 | const tokens = new Set(await Promise.all(Array.from({ length: 20 }, () => store.putAction({ kind: "approval", approvalId: "approval-a" }, owner)))); |
| 114 | assert.equal(tokens.size, 20); |
| 115 | for (const token of tokens) assert.match(token, /^[a-f0-9]{32}$/); |
| 116 | const token = [...tokens][0]; |
| 117 | for (const wrong of [null, { chatId }, { ...owner, chatId: "other" }, { ...owner, threadId: "other" }, |
| 118 | { ...owner, actorId: "" }, { ...owner, actorId: "telegram-user:bob" }, { ...owner, turnId: "other" }]) { |
| 119 | assert.equal(await store.takeAction(token, wrong), null); |
| 120 | } |
| 121 | assert.equal((await store.takeAction(token, owner)).approvalId, "approval-a"); |
| 122 | assert.equal(await store.getAction(token, owner), null); |
| 123 | await assert.rejects(store.putAction({ kind: "approval" }, { chatId, threadId }), ApprovalOwnershipError); |
| 124 | }); |
| 125 | }); |
| 126 | |
| 127 | test("legacy privileged tokens missing thread, turn or human refuse after reload", async () => { |
| 128 | await withStore(async (store) => { |
| 129 | await store.recordTurnOrigin(chatId, threadId, turnId, owner.actorId); |
| 130 | for (const binding of [undefined, { chatId }, { chatId, threadId }, { chatId, threadId, turnId }]) { |
| 131 | const token = String(Object.keys(store.data.actions).length); |
| 132 | store.data.actions[token] = { kind: "approval", owner: binding, approvalId: "approval-a", createdAt: new Date().toISOString() }; |
| 133 | } |
| 134 | await store.save(); |
| 135 | const reopened = await ThreadStore.open(store.filePath, { actions: true }); |
| 136 | for (const token of Object.keys(reopened.data.actions)) assert.equal(await reopened.getAction(token, owner), null); |
| 137 | }); |
| 138 | }); |
| 139 | |
| 140 | test("origins are immutable, concurrent, bounded and retired tokens never acquire newer humans", async () => { |
| 141 | await withStore(async (store) => { |
| 142 | await Promise.all([ |
| 143 | store.recordTurnOrigin(chatId, threadId, turnId, owner.actorId), |
| 144 | store.recordTurnOrigin(chatId, threadId, "turn-b", "telegram-user:bob") |
| 145 | ]); |
| 146 | assert.equal(store.turnOrigin(chatId, threadId, turnId).actorId, owner.actorId); |
| 147 | assert.equal(store.turnOrigin(chatId, threadId, "turn-b").actorId, "telegram-user:bob"); |
| 148 | const token = await store.putAction({ kind: "approval", approvalId: "approval-a" }, owner); |
| 149 | await assert.rejects(store.recordTurnOrigin(chatId, threadId, turnId, "telegram-user:bob"), ApprovalOwnershipError); |
| 150 | await store.recordTurnOrigin(chatId, threadId, "turn-c", "telegram-user:bob"); |
| 151 | const reopened = await ThreadStore.open(store.filePath, { actions: true, actionLimit: 2 }); |
| 152 | assert.equal((await reopened.getChat(chatId)).turnOrigins.length, 2); |
| 153 | assert.equal(reopened.turnOrigin(chatId, threadId, turnId), null); |
| 154 | assert.equal(await reopened.getAction(token, owner), null); |
| 155 | await reopened.setChat(chatId, { threadId: "other" }); |
| 156 | assert.equal(reopened.turnOrigin(chatId, threadId, "turn-c"), null); |
| 157 | await reopened.setChat(chatId, { threadId }); |
| 158 | assert.equal(reopened.turnOrigin(chatId, threadId, "turn-c"), null); |
| 159 | }, { actionLimit: 2 }); |
| 160 | }); |
| 161 | |
| 162 | test("shared decisions require Runtime pending turn and originating human, not current sender", async () => { |
| 163 | await withRuntime(async (runtime) => withStore(async (store) => { |
| 164 | await store.recordTurnOrigin(chatId, threadId, turnId, owner.actorId); |
| 165 | const args = { store, chatId, actorId: owner.actorId, approvalId: "approval-a", decision: "allow" }; |
| 166 | for (const override of [{ actorId: "" }, { actorId: "telegram-user:bob" }, { approvalId: "approval-b" }, |
| 167 | { chatId: "other" }, { turnId: "turn-b" }]) { |
| 168 | await assert.rejects(decideApproval(runtime.runtimeJson, { ...args, ...override }), ApprovalOwnershipError); |
| 169 | } |
| 170 | runtime.state.pending = [{ id: "approval-a" }]; |
| 171 | await assert.rejects(decideApproval(runtime.runtimeJson, args), ApprovalOwnershipError); |
| 172 | runtime.state.pending = [{ id: "approval-a", turn_id: "turn-b" }]; |
| 173 | await assert.rejects(decideApproval(runtime.runtimeJson, args), ApprovalOwnershipError); |
| 174 | assert.equal(runtime.state.posts.length, 0); |
| 175 | runtime.state.pending = [{ id: "approval-a", turn_id: turnId }]; |
| 176 | await decideApproval(runtime.runtimeJson, { ...args, decision: "deny", remember: "yes" }); |
| 177 | assert.deepEqual(runtime.state.posts, [{ route: "/v1/approvals/approval-a", body: { decision: "deny", remember: false } }]); |
| 178 | })); |
| 179 | }); |
| 180 | |
| 181 | test("thread switch during Runtime read refuses before the approval POST", async () => { |
| 182 | await withStore(async (store) => { |
| 183 | await store.recordTurnOrigin(chatId, threadId, turnId, owner.actorId); |
| 184 | const runtimeJson = async (route, options) => { |
| 185 | assert.equal(options, undefined, "must not post after switch"); |
| 186 | await store.setChat(chatId, { threadId: "other" }); |
| 187 | return { pending_approvals: [{ id: "approval-a", turn_id: turnId }] }; |
| 188 | }; |
| 189 | await assert.rejects(decideApproval(runtimeJson, { store, chatId, actorId: owner.actorId, approvalId: "approval-a", decision: "allow" }), ApprovalOwnershipError); |
| 190 | }); |
| 191 | }); |
| 192 | |
| 193 | for (const bridge of Object.keys(identities)) { |
| 194 | test(`${bridge} shipping decision refuses unrelated, missing and other originating human over Runtime HTTP`, async () => { |
| 195 | await withRuntime(async (runtime) => withStore(async (store) => { |
| 196 | await store.recordTurnOrigin(chatId, threadId, turnId, actors[bridge]); |
| 197 | await store.patchChat(chatId, { authorizedIdentity: identities[bridge]("bob"), activeTurnId: "later-turn" }); |
| 198 | const reopened = await ThreadStore.open(store.filePath, { actions: true }); |
| 199 | const { decideApproval: consumer } = await shipping(bridge, ["turnActor", "decideApproval"], environment(reopened, runtime, |
| 200 | bridge === "weixin-bridge" ? { botAccount: { accountId: "bot-a" } } : {})); |
| 201 | const action = { decision: "allow", approvalId: "approval-a", remember: true }; |
| 202 | await consumer(chatId, { ...action, approvalId: "approval-b" }, identities[bridge]("alice")); |
| 203 | if (bridge !== "weixin-bridge") { |
| 204 | await consumer(chatId, action, identities[bridge]("bob")); |
| 205 | await consumer(chatId, action, identities[bridge]("")); |
| 206 | } |
| 207 | assert.equal(runtime.state.posts.length, 0); |
| 208 | await consumer(chatId, action, identities[bridge]("alice")); |
| 209 | assert.deepEqual(runtime.state.posts, [{ route: "/v1/approvals/approval-a", body: { decision: "allow", remember: bridge !== "wecom-bridge" } }]); |
| 210 | })); |
| 211 | }); |
| 212 | } |
| 213 | |
| 214 | test("Weixin shipping decision rejects replacement or missing account identity without transferring the origin", async () => { |
| 215 | await withRuntime(async (runtime) => withStore(async (store) => { |
| 216 | await store.recordTurnOrigin(chatId, threadId, turnId, actors["weixin-bridge"]); |
| 217 | const reopened = await ThreadStore.open(store.filePath, { actions: true }); |
| 218 | const botAccount = { accountId: "bot-a" }; |
| 219 | const { decideApproval: consumer } = await shipping("weixin-bridge", ["turnActor", "decideApproval"], |
| 220 | environment(reopened, runtime, { botAccount })); |
| 221 | const action = { decision: "allow", approvalId: "approval-a", remember: true }; |
| 222 | for (const accountId of ["bot-b", "", null]) { |
| 223 | botAccount.accountId = accountId; |
| 224 | await consumer(chatId, action); |
| 225 | assert.equal(runtime.state.posts.length, 0); |
| 226 | assert.equal(reopened.turnOrigin(chatId, threadId, turnId).actorId, actors["weixin-bridge"]); |
| 227 | } |
| 228 | botAccount.accountId = "bot-a"; |
| 229 | await consumer(chatId, action); |
| 230 | assert.deepEqual(runtime.state.posts, [{ route: "/v1/approvals/approval-a", body: { decision: "allow", remember: true } }]); |
| 231 | })); |
| 232 | }); |
| 233 | |
| 234 | test("Weixin shipping decision refuses a legacy origin without bot account provenance after reload", async () => { |
| 235 | await withRuntime(async (runtime) => withStore(async (store) => { |
| 236 | await store.recordTurnOrigin(chatId, threadId, turnId, "weixin-user:chat-a"); |
| 237 | const reopened = await ThreadStore.open(store.filePath, { actions: true }); |
| 238 | const { decideApproval: consumer } = await shipping("weixin-bridge", ["turnActor", "decideApproval"], |
| 239 | environment(reopened, runtime, { botAccount: { accountId: "bot-a" } })); |
| 240 | await consumer(chatId, { decision: "allow", approvalId: "approval-a", remember: true }); |
| 241 | assert.equal(runtime.state.posts.length, 0); |
| 242 | assert.equal(reopened.turnOrigin(chatId, threadId, turnId).actorId, "weixin-user:chat-a"); |
| 243 | })); |
| 244 | }); |
| 245 | |
| 246 | for (const bridge of ["telegram-bridge", "feishu-bridge", "wecom-bridge"]) { |
| 247 | test(`${bridge} accepted start binds before delivery; failed, queued and later starts cannot transfer an old approval`, async () => { |
| 248 | await withRuntime(async (runtime) => withStore(async (store) => { |
| 249 | const functions = await shipping(bridge, ["turnActor", "runPrompt", "decideApproval"], environment(store, runtime, bridge === "wecom-bridge" ? { |
| 250 | streamTurnEvents: async (chat, frame, thread, turn) => { |
| 251 | const reopened = await ThreadStore.open(store.filePath, { actions: true }); |
| 252 | assert.ok(reopened.turnOrigin(chat, thread, turn), "origin must be durable before delivery"); |
| 253 | } |
| 254 | } : {})); |
| 255 | await invokeStart(bridge, functions.runPrompt, "alice"); |
| 256 | const first = runtime.state.accepted[0]; |
| 257 | assert.equal(store.turnOrigin(chatId, threadId, first).actorId, actors[bridge]); |
| 258 | runtime.state.pending = [{ id: "approval-a", turn_id: first }]; |
| 259 | await store.patchChat(chatId, { authorizedIdentity: identities[bridge]("bob") }); |
| 260 | runtime.state.failStart = true; |
| 261 | await assert.rejects(invokeStart(bridge, functions.runPrompt, "bob"), /start rejected/); |
| 262 | assert.equal(store.turnOrigin(chatId, threadId, first).actorId, actors[bridge]); |
| 263 | runtime.state.failStart = false; |
| 264 | runtime.state.turns = [{ id: first, status: "queued" }]; |
| 265 | const before = runtime.state.startAttempts; |
| 266 | await invokeStart(bridge, functions.runPrompt, "bob"); |
| 267 | assert.equal(runtime.state.startAttempts, before); |
| 268 | runtime.state.turns = []; |
| 269 | await invokeStart(bridge, functions.runPrompt, "bob"); |
| 270 | assert.equal(store.turnOrigin(chatId, threadId, first).actorId, actors[bridge]); |
| 271 | const reopened = await ThreadStore.open(store.filePath, { actions: true }); |
| 272 | const { decideApproval: consumer } = await shipping(bridge, ["turnActor", "decideApproval"], environment(reopened, runtime)); |
| 273 | await consumer(chatId, { decision: "allow", approvalId: "approval-a" }, identities[bridge]("bob")); |
| 274 | assert.equal(runtime.state.posts.length, 0); |
| 275 | await consumer(chatId, { decision: "allow", approvalId: "approval-a" }, identities[bridge]("alice")); |
| 276 | assert.equal(runtime.state.posts.length, 1); |
| 277 | })); |
| 278 | }); |
| 279 | } |
| 280 | |
| 281 | test("Telegram shipping SSE mints a bound button and actual callback human controls retry, restart and replay", async () => { |
| 282 | await withRuntime(async (runtime) => withStore(async (store) => { |
| 283 | await store.recordTurnOrigin(chatId, threadId, turnId, owner.actorId); |
| 284 | runtime.state.events = [{ seq: 1, turn_id: turnId, event: "approval.required", payload: { approval_id: "approval-a" } }]; |
| 285 | const telegram = await import("../../telegram-bridge/src/lib.mjs"); |
| 286 | const env = environment(store, runtime, { |
| 287 | authHeaders: runtime.authHeaders, readSse, readJsonSafe, compactRuntimeError, |
| 288 | sendTypingAction: async () => {}, TYPING_INTERVAL_MS: 2000, LAST_SEQ_FLUSH_INTERVAL_MS: 2000, |
| 289 | approvalKeyboard: telegram.approvalKeyboard |
| 290 | }); |
| 291 | const { streamTurnEvents } = await shipping("telegram-bridge", ["streamTurnEvents"], env); |
| 292 | await streamTurnEvents(chatId, threadId, turnId, 0); |
| 293 | const [token] = Object.keys(store.data.actions); |
| 294 | assert.ok(token, "actual SSE must mint a token"); |
| 295 | assert.deepEqual(store.data.actions[token].owner, owner); |
| 296 | const invoke = async (current, user) => { |
| 297 | const values = environment(current, runtime, { |
| 298 | ...telegram, config: { ...env.config, allowGroups: true, allowUnlisted: true, allowlist: [] }, |
| 299 | answerCallback: async () => {}, editMessageReplyMarkup: async () => {}, |
| 300 | rememberAuthorizedIdentity: async (identity) => current.patchChat(chatId, { authorizedIdentity: identity }) |
| 301 | }); |
| 302 | const { handleCallbackQuery } = await shipping("telegram-bridge", |
| 303 | ["turnActor", "decideApproval", "handleStoredAction", "handleModalAction", "handleCallbackQuery"], values); |
| 304 | return handleCallbackQuery({ id: "callback-fixture", data: telegram.approvalKeyboard(token).inline_keyboard[0][0].callback_data, |
| 305 | from: { id: user }, message: { message_id: 1, chat: { id: chatId, type: "group" } } }); |
| 306 | }; |
| 307 | await invoke(store, "bob"); |
| 308 | assert.equal(runtime.state.posts.length, 0); |
| 309 | // The same approval id cannot authorize a different accepted turn, even |
| 310 | // when the same human happens to have initiated both turns. |
| 311 | await store.recordTurnOrigin(chatId, threadId, "other-turn", owner.actorId); |
| 312 | runtime.state.pending = [{ id: "approval-a", turn_id: "other-turn" }]; |
| 313 | await invoke(store, "alice"); |
| 314 | assert.equal(runtime.state.posts.length, 0); |
| 315 | runtime.state.pending = [{ id: "approval-a", turn_id: turnId }]; |
| 316 | runtime.state.acceptDecision = false; |
| 317 | await assert.rejects(invoke(store, "alice"), /fixture unavailable/); |
| 318 | const reopened = await ThreadStore.open(store.filePath, { actions: true }); |
| 319 | assert.ok(await reopened.getAction(token, owner)); |
| 320 | runtime.state.acceptDecision = true; |
| 321 | await invoke(reopened, "alice"); |
| 322 | const accepted = await ThreadStore.open(store.filePath, { actions: true }); |
| 323 | assert.equal(await accepted.getAction(token, owner), null); |
| 324 | await invoke(accepted, "alice"); |
| 325 | assert.equal(runtime.state.posts.length, 2, "failed then accepted; accepted replay never posts"); |
| 326 | })); |
| 327 | }); |
| 328 | |
| 329 | test("Telegram and Feishu actual dispatch passes the admitted human to prompt and text approval", async () => { |
| 330 | for (const bridge of ["telegram-bridge", "feishu-bridge"]) { |
| 331 | const lib = await import(`../../${bridge}/src/lib.mjs`); |
| 332 | const calls = []; |
| 333 | const identity = identities[bridge]("alice"); |
| 334 | const values = { ...lib, sendText: async () => {}, |
| 335 | decideApproval: async (...args) => calls.push(args), |
| 336 | startPromptTurn: (...args) => calls.push(args), |
| 337 | runPrompt: async (...args) => calls.push(args) }; |
| 338 | const { handleCommand } = await shipping(bridge, ["handleCommand"], values); |
| 339 | await handleCommand(chatId, lib.parseCommand("/allow approval-a"), identity); |
| 340 | await handleCommand(chatId, lib.parseCommand("hello"), identity); |
| 341 | assert.equal(calls.length, 2, bridge); |
| 342 | assert.equal(calls[0][2], identity, bridge); |
| 343 | assert.equal(calls[1][2], identity, bridge); |
| 344 | } |
| 345 | }); |
| 346 | |
| 347 | |
| 348 | test("Telegram and Feishu shipping inbound events retain the actual admitted human through accepted start", async () => { |
| 349 | for (const bridge of ["telegram-bridge", "feishu-bridge"]) { |
| 350 | await withRuntime(async (runtime) => withStore(async (store) => { |
| 351 | const lib = await import(`../../${bridge}/src/lib.mjs`); |
| 352 | let delivered; |
| 353 | const delivery = new Promise((resolve) => { delivered = resolve; }); |
| 354 | const values = environment(store, runtime, { |
| 355 | ...lib, activeTurnTasks: new Map(), |
| 356 | config: { allowGroups: true, allowUnlisted: false, allowlist: [chatId], requirePrefixInGroup: false }, |
| 357 | streamTurnEvents: async (chat, thread, turn) => { |
| 358 | const reopened = await ThreadStore.open(store.filePath, { actions: true }); |
| 359 | assert.equal(reopened.turnOrigin(chat, thread, turn).actorId, actors[bridge]); |
| 360 | delivered(); |
| 361 | } |
| 362 | }); |
| 363 | const names = bridge === "telegram-bridge" ? |
| 364 | ["turnActor", "runPrompt", "startPromptTurn", "handleCommand", "rememberAuthorizedIdentity", "handleIncomingUpdate"] : |
| 365 | ["turnActor", "runPrompt", "handleCommand", "handleIncomingMessage"]; |
| 366 | const handlers = await shipping(bridge, names, values); |
| 367 | if (bridge === "telegram-bridge") { |
| 368 | await handlers.handleIncomingUpdate({ update_id: 1, message: { message_id: 1, text: "hello", |
| 369 | chat: { id: chatId, type: "group" }, from: { id: "alice" } } }); |
| 370 | await delivery; |
| 371 | } else { |
| 372 | await handlers.handleIncomingMessage({ sender: { sender_id: { open_id: "alice" } }, |
| 373 | message: { chat_id: chatId, chat_type: "group", message_id: "message-a", message_type: "text", |
| 374 | content: JSON.stringify({ text: "hello" }) } }); |
| 375 | } |
| 376 | assert.equal(runtime.state.accepted.length, 1); |
| 377 | assert.equal(store.turnOrigin(chatId, threadId, runtime.state.accepted[0]).actorId, actors[bridge]); |
| 378 | })); |
| 379 | } |
| 380 | }); |
| 381 | |
| 382 | test("WeCom shipping SSE and natural approval callback use durable Runtime-turn origin, not the hint's latest sender", async () => { |
| 383 | await withRuntime(async (runtime) => withStore(async (store) => { |
| 384 | const lib = await import("../../wecom-bridge/src/lib.mjs"); |
| 385 | await store.recordTurnOrigin(chatId, threadId, turnId, actors["wecom-bridge"]); |
| 386 | runtime.state.events = [{ seq: 1, turn_id: turnId, event: "approval.required", payload: { approval_id: "approval-a" } }]; |
| 387 | const pendingApprovals = new Map(); |
| 388 | const values = environment(store, runtime, { |
| 389 | ...lib, authHeaders: runtime.authHeaders, readSse, readJsonSafe, compactRuntimeError, |
| 390 | config: { runtimeUrl: runtime.runtimeUrl, turnTimeoutMs: 1000, approvalTimeoutMs: 10000 }, |
| 391 | pendingApprovals, generateReqId: () => "stream-fixture", client: { replyStream: async () => {} } |
| 392 | }); |
| 393 | const { streamTurnEvents } = await shipping("wecom-bridge", ["streamTurnEvents"], values); |
| 394 | await streamTurnEvents(chatId, identities["wecom-bridge"]("bob"), threadId, turnId, 0); |
| 395 | assert.equal(pendingApprovals.get(chatId).approvalId, "approval-a"); |
| 396 | const reopened = await ThreadStore.open(store.filePath, { actions: true }); |
| 397 | const { handleCommand } = await shipping("wecom-bridge", ["turnActor", "decideApproval", "handleCommand"], |
| 398 | { ...values, threadStore: reopened }); |
| 399 | await handleCommand(chatId, lib.parseCommand("允许"), identities["wecom-bridge"]("bob")); |
| 400 | assert.equal(runtime.state.posts.length, 0); |
| 401 | assert.ok(pendingApprovals.has(chatId), "refused member must not retire the legitimate pending hint"); |
| 402 | await handleCommand(chatId, lib.parseCommand("允许"), identities["wecom-bridge"]("alice")); |
| 403 | assert.equal(runtime.state.posts.length, 1); |
| 404 | assert.deepEqual(runtime.state.posts[0].body, { decision: "allow", remember: false }); |
| 405 | assert.equal(pendingApprovals.has(chatId), false); |
| 406 | })); |
| 407 | }); |
| 408 | |
| 409 | |
| 410 | test("retired origin refuses stale pending approval even when newer turns have the same human", async () => { |
| 411 | await withRuntime(async (runtime) => withStore(async (store) => { |
| 412 | await store.recordTurnOrigin(chatId, threadId, turnId, owner.actorId); |
| 413 | const token = await store.putAction({ kind: "approval", approvalId: "approval-a" }, owner); |
| 414 | await store.recordTurnOrigin(chatId, threadId, "newer-1", owner.actorId); |
| 415 | await store.recordTurnOrigin(chatId, threadId, "newer-2", owner.actorId); |
| 416 | const reopened = await ThreadStore.open(store.filePath, { actions: true, actionLimit: 2 }); |
| 417 | assert.equal(await reopened.getAction(token, owner), null); |
| 418 | await assert.rejects(decideApproval(runtime.runtimeJson, { store: reopened, chatId, actorId: owner.actorId, |
| 419 | approvalId: "approval-a", decision: "allow" }), ApprovalOwnershipError); |
| 420 | assert.equal(runtime.state.posts.length, 0); |
| 421 | }, { actionLimit: 2 })); |
| 422 | }); |
| 423 | |
| 424 | test("existing action TTL is enforced on lookup after restart", async () => { |
| 425 | await withStore(async (store) => { |
| 426 | await store.recordTurnOrigin(chatId, threadId, turnId, owner.actorId); |
| 427 | const token = await store.putAction({ kind: "approval", approvalId: "approval-a" }, owner); |
| 428 | store.data.actions[token].createdAt = new Date(0).toISOString(); |
| 429 | await store.save(); |
| 430 | const reopened = await ThreadStore.open(store.filePath, { actions: true }); |
| 431 | assert.equal(await reopened.getAction(token, owner), null); |
| 432 | }); |
| 433 | }); |
| 434 |