| 1 | import test from "node:test"; |
| 2 | import assert from "node:assert/strict"; |
| 3 | import { mkdtemp, rm } from "node:fs/promises"; |
| 4 | import { getEventListeners } from "node:events"; |
| 5 | import { createServer } from "node:http"; |
| 6 | import { tmpdir } from "node:os"; |
| 7 | import path from "node:path"; |
| 8 | |
| 9 | import { ThreadStore } from "../../bridge-core/src/lib.mjs"; |
| 10 | |
| 11 | import { |
| 12 | activeTurnBlock, |
| 13 | apiPost, |
| 14 | commandAction, |
| 15 | extractText, |
| 16 | MessageItemType, |
| 17 | parseBool, |
| 18 | parseCommand, |
| 19 | parseList, |
| 20 | preservedChatStateFields, |
| 21 | processUpdateBatch, |
| 22 | sendMessage, |
| 23 | splitMessage |
| 24 | } from "../src/lib.mjs"; |
| 25 | |
| 26 | test("extractText reads text and voice transcript items", () => { |
| 27 | assert.equal( |
| 28 | extractText([{ type: MessageItemType.TEXT, text_item: { text: "hello" } }]), |
| 29 | "hello" |
| 30 | ); |
| 31 | assert.equal( |
| 32 | extractText([{ type: MessageItemType.VOICE, voice_item: { text: "voice text" } }]), |
| 33 | "voice text" |
| 34 | ); |
| 35 | }); |
| 36 | |
| 37 | test("shared command helpers preserve Weixin bridge command behavior", () => { |
| 38 | assert.deepEqual(parseList("u1, u2 ,, "), ["u1", "u2"]); |
| 39 | assert.equal(parseBool("yes"), true); |
| 40 | assert.deepEqual(parseCommand("/allow ap_1 remember"), { |
| 41 | name: "allow", |
| 42 | args: "ap_1 remember" |
| 43 | }); |
| 44 | assert.deepEqual(commandAction(parseCommand("/model auto")), { |
| 45 | kind: "set_model", |
| 46 | modelName: "auto" |
| 47 | }); |
| 48 | assert.deepEqual(commandAction(parseCommand("/unknown value")), { |
| 49 | kind: "prompt", |
| 50 | prompt: "/unknown value" |
| 51 | }); |
| 52 | }); |
| 53 | |
| 54 | test("shared state and runtime helpers preserve Weixin bridge behavior", () => { |
| 55 | assert.deepEqual(preservedChatStateFields({ model: "m", activeTurnId: "turn-1" }), { |
| 56 | model: "m" |
| 57 | }); |
| 58 | assert.deepEqual(splitMessage("a🧪b", 2), ["a🧪", "b"]); |
| 59 | assert.deepEqual(activeTurnBlock({ turns: [{ id: "turn-1", status: "in_progress" }] }), { |
| 60 | turnId: "turn-1", |
| 61 | message: "Thread already has active turn turn-1. Wait for it to finish or send /interrupt." |
| 62 | }); |
| 63 | }); |
| 64 | |
| 65 | test("a crash mid-batch replays the batch without losing or re-running prompts", async () => { |
| 66 | const dir = await mkdtemp(path.join(tmpdir(), "codewhale-weixin-bridge-")); |
| 67 | try { |
| 68 | const statePath = path.join(dir, "thread-map.json"); |
| 69 | const batch = [1, 2, 3].map((id) => ({ from_user_id: "u", message_id: id })); |
| 70 | const keyOf = (msg) => `${msg.from_user_id}:${msg.message_id}`; |
| 71 | const handled = []; |
| 72 | const cursors = []; |
| 73 | const crash = new Error("process killed"); |
| 74 | const first = await ThreadStore.open(statePath, { messageLimit: 50 }); |
| 75 | let afterCrash = null; |
| 76 | await assert.rejects(processUpdateBatch({ |
| 77 | messages: batch, nextCursor: "buf-2", store: first, keyOf, |
| 78 | handle: async (msg) => { |
| 79 | if (msg.message_id === 2) { |
| 80 | // What a restarted process finds on disk while message 2 runs. |
| 81 | afterCrash = await ThreadStore.open(statePath, { messageLimit: 50 }); |
| 82 | throw crash; |
| 83 | } |
| 84 | handled.push(msg.message_id); |
| 85 | }, |
| 86 | interrupted: async () => assert.fail("nothing is interrupted on the first pass"), |
| 87 | commitCursor: async (buf) => cursors.push(buf), |
| 88 | }), crash); |
| 89 | assert.deepEqual(cursors, [], "the cursor does not move past unhandled messages"); |
| 90 | |
| 91 | const reported = []; |
| 92 | await processUpdateBatch({ |
| 93 | messages: batch, nextCursor: "buf-2", store: afterCrash, keyOf, |
| 94 | handle: async (msg) => handled.push(msg.message_id), |
| 95 | interrupted: async (msg) => reported.push(msg.message_id), |
| 96 | commitCursor: async (buf) => cursors.push(buf), |
| 97 | }); |
| 98 | assert.deepEqual(handled, [1, 3], "message 1 is not re-run; message 3 is not lost"); |
| 99 | assert.deepEqual(reported, [2], "the in-flight message is reported, not re-run"); |
| 100 | assert.deepEqual(cursors, ["buf-2"]); |
| 101 | } finally { |
| 102 | await rm(dir, { recursive: true, force: true }); |
| 103 | } |
| 104 | }); |
| 105 | |
| 106 | test("sendMessage checks application responses from a local HTTP server", async (t) => { |
| 107 | let responseBody; |
| 108 | let status = 200; |
| 109 | const requests = []; |
| 110 | const server = createServer((req, res) => { |
| 111 | requests.push({ method: req.method, url: req.url }); |
| 112 | req.resume(); |
| 113 | res.writeHead(status, { "Content-Type": "application/json", Connection: "close" }); |
| 114 | res.end(responseBody); |
| 115 | }); |
| 116 | await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); |
| 117 | t.after(() => new Promise((resolve) => server.close(resolve))); |
| 118 | const params = { |
| 119 | baseUrl: `http://127.0.0.1:${server.address().port}`, |
| 120 | token: "test-token", |
| 121 | body: { msg: { to_user_id: "test-user", context_token: "test-context" } }, |
| 122 | }; |
| 123 | const cases = [ |
| 124 | { name: "ret zero is API acceptance", body: '{"ret":0}', accepted: { ret: 0 } }, |
| 125 | { name: "errcode zero is API acceptance", body: '{"errcode":0}', accepted: { errcode: 0 } }, |
| 126 | { name: "omitted status matches Tencent reference compatibility", body: '{}', accepted: {} }, |
| 127 | { name: "HTTP 200 with nonzero ret is rejected", body: '{"ret":-1,"errmsg":"response-secret"}', outcome: "rejected", message: /ret=-1/ }, |
| 128 | { name: "HTTP 200 with nonzero errcode is rejected", body: '{"ret":0,"errcode":-14,"errmsg":"response-secret"}', outcome: "rejected", message: /errcode=-14/ }, |
| 129 | { name: "malformed JSON is uncertain", body: '<html>response-secret</html>', outcome: "uncertain", message: /invalid JSON/ }, |
| 130 | { name: "empty response is uncertain", body: '', outcome: "uncertain", message: /invalid JSON/ }, |
| 131 | { name: "null response is uncertain", body: 'null', outcome: "uncertain", message: /non-object/ }, |
| 132 | { name: "array response is uncertain", body: '[]', outcome: "uncertain", message: /non-object/ }, |
| 133 | { name: "string response is uncertain", body: '"response-secret"', outcome: "uncertain", message: /non-object/ }, |
| 134 | { name: "string status is uncertain", body: '{"ret":"response-secret"}', outcome: "uncertain", message: /invalid status/ }, |
| 135 | { name: "null status is uncertain", body: '{"ret":null}', outcome: "uncertain", message: /invalid status/ }, |
| 136 | { name: "boolean errcode is uncertain", body: '{"ret":0,"errcode":false}', outcome: "uncertain", message: /invalid status/ }, |
| 137 | { name: "HTTP failure is uncertain and does not expose its body", status: 503, body: 'response-secret', outcome: "uncertain", message: /HTTP 503/ }, |
| 138 | ]; |
| 139 | for (const row of cases) { |
| 140 | await t.test(row.name, async () => { |
| 141 | responseBody = row.body; |
| 142 | status = row.status ?? 200; |
| 143 | if (Object.hasOwn(row, "accepted")) { |
| 144 | assert.deepEqual(await sendMessage(params), row.accepted); |
| 145 | } else { |
| 146 | await assert.rejects(sendMessage(params), (error) => { |
| 147 | assert.equal(error.deliveryStatus, row.outcome); |
| 148 | assert.match(error.message, row.message); |
| 149 | assert.doesNotMatch(error.message, /response-secret/); |
| 150 | return true; |
| 151 | }); |
| 152 | } |
| 153 | }); |
| 154 | } |
| 155 | assert.equal(requests.length, cases.length); |
| 156 | assert.ok(requests.every((req) => req.method === "POST" && req.url === "/ilink/bot/sendmessage")); |
| 157 | }); |
| 158 | |
| 159 | test("sendMessage preserves network errors and safely handles non-Error failures", async (t) => { |
| 160 | let failure = new Error("network unavailable"); |
| 161 | t.mock.method(globalThis, "fetch", async () => { throw failure; }); |
| 162 | const params = { baseUrl: "http://localhost", body: { msg: {} } }; |
| 163 | await assert.rejects(sendMessage(params), (error) => error === failure && error.deliveryStatus === "uncertain"); |
| 164 | failure = "response-secret"; |
| 165 | await assert.rejects(sendMessage(params), (error) => { |
| 166 | assert.equal(error.deliveryStatus, "uncertain"); |
| 167 | assert.equal(error.message, "iLink sendMessage transport failed"); |
| 168 | return true; |
| 169 | }); |
| 170 | }); |
| 171 | |
| 172 | test("apiPost does not start a request for an already-aborted signal", async (t) => { |
| 173 | const controller = new AbortController(); |
| 174 | const reason = new Error("bridge stopped"); |
| 175 | controller.abort(reason); |
| 176 | t.mock.method(globalThis, "fetch", () => assert.fail("aborted request must not reach fetch")); |
| 177 | await assert.rejects(apiPost({ baseUrl: "http://localhost", endpoint: "test", body: "{}", signal: controller.signal }), reason); |
| 178 | assert.equal(getEventListeners(controller.signal, "abort").length, 0); |
| 179 | }); |
| 180 | |
| 181 | test("apiPost removes external abort listeners after success and failure", async (t) => { |
| 182 | const controller = new AbortController(); |
| 183 | let fail = false; |
| 184 | t.mock.method(globalThis, "fetch", async () => { |
| 185 | assert.equal(getEventListeners(controller.signal, "abort").length, 1); |
| 186 | if (fail) throw new Error("network unavailable"); |
| 187 | return new Response("{}"); |
| 188 | }); |
| 189 | const params = { baseUrl: "http://localhost", endpoint: "test", body: "{}", signal: controller.signal }; |
| 190 | assert.equal(await apiPost(params), "{}"); |
| 191 | assert.equal(getEventListeners(controller.signal, "abort").length, 0); |
| 192 | fail = true; |
| 193 | await assert.rejects(apiPost(params), /network unavailable/); |
| 194 | assert.equal(getEventListeners(controller.signal, "abort").length, 0); |
| 195 | }); |
| 196 | |
| 197 | test("apiPost propagates in-flight cancellation and clears its external listener", async (t) => { |
| 198 | const controller = new AbortController(); |
| 199 | let requestSignal; |
| 200 | t.mock.method(globalThis, "fetch", (_url, options) => { |
| 201 | requestSignal = options.signal; |
| 202 | return new Promise((_resolve, reject) => { |
| 203 | requestSignal.addEventListener("abort", () => reject(requestSignal.reason), { once: true }); |
| 204 | }); |
| 205 | }); |
| 206 | const pending = apiPost({ baseUrl: "http://localhost", endpoint: "test", body: "{}", signal: controller.signal }); |
| 207 | const reason = new Error("bridge stopped"); |
| 208 | controller.abort(reason); |
| 209 | await assert.rejects(pending, reason); |
| 210 | assert.equal(requestSignal.aborted, true); |
| 211 | assert.equal(getEventListeners(controller.signal, "abort").length, 0); |
| 212 | }); |
| 213 |