| 1 | import test from "node:test"; |
| 2 | import assert from "node:assert/strict"; |
| 3 | import http from "node:http"; |
| 4 | import fs from "node:fs/promises"; |
| 5 | import os from "node:os"; |
| 6 | import path from "node:path"; |
| 7 | import { spawn } from "node:child_process"; |
| 8 | import { fileURLToPath } from "node:url"; |
| 9 | |
| 10 | const entry = fileURLToPath(new URL("../src/index.mjs", import.meta.url)); |
| 11 | const delay = (ms) => new Promise((resolve) => setTimeout(resolve, ms)); |
| 12 | async function until(check, message, timeout = 8000) { |
| 13 | const deadline = Date.now() + timeout; |
| 14 | while (Date.now() < deadline) { if (await check()) return; await delay(25); } |
| 15 | assert.fail(message); |
| 16 | } |
| 17 | function incoming(id, text, user = "alice") { |
| 18 | return { message_id: id, from_user_id: user, message_type: 1, context_token: "context-fixture", item_list: [{ type: 1, text_item: { text } }] }; |
| 19 | } |
| 20 | async function fixture(t, options = {}) { |
| 21 | const dir = await fs.mkdtemp(path.join(os.tmpdir(), "weixin-runtime-")); |
| 22 | const state = { polls: 0, batches: [], posts: [], sent: [], lookups: 0, streams: new Set(), turns: [], items: [], seq: 0, approvals: [], decisions: [], logs: "" }; |
| 23 | const json = (response, body, status = 200) => { response.writeHead(status, { "Content-Type": "application/json" }); response.end(JSON.stringify(body)); }; |
| 24 | const server = http.createServer(async (req, res) => { |
| 25 | let raw = ""; for await (const chunk of req) raw += chunk; |
| 26 | const body = raw ? JSON.parse(raw) : {}; |
| 27 | const url = new URL(req.url, "http://fixture"); |
| 28 | if (url.pathname.endsWith("/getupdates")) { |
| 29 | state.polls++; |
| 30 | await delay(35); |
| 31 | if (res.destroyed) return; |
| 32 | return json(res, { ret: 0, msgs: state.batches.shift() || [], get_updates_buf: String(state.polls) }); |
| 33 | } |
| 34 | if (url.pathname.endsWith("/sendmessage")) { |
| 35 | state.sent.push(body.msg); |
| 36 | if (options.send) return options.send(res, body.msg, state, json); |
| 37 | return json(res, { ret: 0 }); |
| 38 | } |
| 39 | if (url.pathname.includes("/msg/notify")) return json(res, { ret: 0 }); |
| 40 | if (url.pathname === "/v1/runtime/info") return json(res, { version: "fixture", capabilities: options.capabilities ?? { turn_operation_idempotency: true, turn_operation_lookup: true } }); |
| 41 | if (url.pathname === "/health") return json(res, { status: "ok" }); |
| 42 | if (url.pathname === "/v1/workspace/status") return json(res, { workspace: dir }); |
| 43 | if (url.pathname === "/v1/threads" && req.method === "POST") { |
| 44 | state.threadRequest = body; |
| 45 | return json(res, { id: "thread-1" }); |
| 46 | } |
| 47 | if (url.pathname === "/v1/threads/thread-1") return json(res, { id: "thread-1", latest_seq: state.seq, turns: state.turns, items: state.items, pending_approvals: state.approvals }); |
| 48 | if (url.pathname === "/v1/threads/thread-1/turns") { |
| 49 | state.posts.push(body); |
| 50 | const turn = { id: "turn-1", thread_id: "thread-1", status: "in_progress" }; |
| 51 | state.turns = [turn]; |
| 52 | if (options.admission) return options.admission(req, res, turn, state, json); |
| 53 | return json(res, { turn }, 201); |
| 54 | } |
| 55 | if (url.pathname.includes("/turn-operations/")) { |
| 56 | state.lookups++; |
| 57 | if (options.lookup) return options.lookup(res, state, json); |
| 58 | return json(res, state.turns[0] || {}, state.turns.length ? 200 : 404); |
| 59 | } |
| 60 | if (url.pathname === "/v1/threads/thread-1/events") { |
| 61 | res.writeHead(200, { "Content-Type": "text/event-stream" }); res.flushHeaders(); |
| 62 | state.streams.add(res); res.on("close", () => state.streams.delete(res)); |
| 63 | return; |
| 64 | } |
| 65 | if (url.pathname === "/v1/approvals/approval-1") { state.decisions.push(body); state.approvals = []; return json(res, {}); } |
| 66 | return json(res, { error: url.pathname }, 404); |
| 67 | }); |
| 68 | await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); |
| 69 | const baseUrl = `http://127.0.0.1:${server.address().port}`; |
| 70 | const account = { accountId: "bot-A", token: "fixture-bot-token", baseUrl }; |
| 71 | await fs.writeFile(path.join(dir, "account.json"), JSON.stringify(account)); |
| 72 | const children = []; |
| 73 | state.start = (extra = {}) => { |
| 74 | const child = spawn(process.execPath, [entry], { env: { ...process.env, CODEWHALE_RUNTIME_URL: baseUrl, CODEWHALE_RUNTIME_TOKEN: "fixture-runtime-token", CODEWHALE_WORKSPACE: dir, WEIXIN_STATE_DIR: dir, WEIXIN_CHAT_ALLOWLIST: "alice,bob", CODEWHALE_TURN_TIMEOUT_MS: "1500", ...extra }, stdio: ["ignore", "pipe", "pipe"] }); |
| 75 | children.push(child); |
| 76 | child.stdout.on("data", (data) => { state.logs += data; }); |
| 77 | child.stderr.on("data", (data) => { state.logs += data; }); |
| 78 | return child; |
| 79 | }; |
| 80 | state.disk = async () => { try { return JSON.parse(await fs.readFile(path.join(dir, "thread-map.json"), "utf8")); } catch { return {}; } }; |
| 81 | state.event = (event, payload, turnId = "turn-1") => { |
| 82 | state.seq++; |
| 83 | const record = { event, payload, turn_id: turnId, seq: state.seq }; |
| 84 | for (const response of state.streams) response.write(`data: ${JSON.stringify(record)}\n\n`); |
| 85 | }; |
| 86 | state.complete = (text, status = "completed") => { |
| 87 | state.items = [{ turn_id: "turn-1", kind: "agent_message", detail: text, summary: "truncated" }]; |
| 88 | state.turns[0].status = status; |
| 89 | state.event("turn.completed", { turn: state.turns[0] }); |
| 90 | }; |
| 91 | state.dir = dir; state.account = account; |
| 92 | t.after(async () => { |
| 93 | for (const child of children) child.kill("SIGKILL"); |
| 94 | for (const response of state.streams) response.destroy(); |
| 95 | server.closeAllConnections(); |
| 96 | await new Promise((resolve) => server.close(resolve)); |
| 97 | await fs.rm(dir, { recursive: true, force: true }); |
| 98 | }); |
| 99 | return state; |
| 100 | } |
| 101 | |
| 102 | test("production process keeps polling and owner approvals live, dedupes real message IDs", async (t) => { |
| 103 | const f = await fixture(t); f.batches.push([incoming(0, "private original prompt"), incoming(undefined, "must not run")]); f.start(); |
| 104 | await until(() => f.streams.size, "accepted turn did not start observation"); |
| 105 | const polls = f.polls; |
| 106 | f.approvals = [{ id: "approval-1", turn_id: "turn-1", tool_name: "fixture-tool" }]; |
| 107 | f.event("approval.required", f.approvals[0]); |
| 108 | await until(() => f.sent.some((msg) => msg.item_list[0].text_item.text.includes("approval_id=approval-1")), "approval notice missing"); |
| 109 | f.batches.push([incoming(0, "private original prompt"), incoming(2, "/allow approval-1", "bob"), incoming(3, "/allow approval-1")]); |
| 110 | await until(() => f.decisions.length === 1, "initiating owner could not approve during stream"); |
| 111 | assert.equal(f.posts.length, 1); assert.ok(f.polls > polls); assert.deepEqual(f.decisions, [{ decision: "allow", remember: false }]); |
| 112 | f.event("item.delta", { kind: "agent_message", delta: "answer" }); f.complete("answer"); |
| 113 | await until(async () => !(await f.disk()).chats?.alice?.activeTurnId, "accepted output did not settle"); |
| 114 | const disk = await f.disk(); assert.ok(disk.messages.length >= 3); assert.equal(disk.chats.alice.contextToken, "context-fixture"); |
| 115 | assert.equal(f.logs.includes("private original prompt"), false); |
| 116 | }); |
| 117 | |
| 118 | test("restart reattaches exact accepted turn and preserves accumulator across EOF", async (t) => { |
| 119 | const f = await fixture(t); f.batches.push([incoming(1, "work")]); const first = f.start(); |
| 120 | await until(() => f.streams.size, "first stream missing"); f.event("item.delta", { kind: "agent_message", delta: "first " }); |
| 121 | await until(async () => (await f.disk()).chats?.alice?.turnDelivery?.responseText === "first ", "delta was not durable"); |
| 122 | for (const response of f.streams) response.end(); |
| 123 | await delay(75); assert.equal((await f.disk()).chats.alice.activeTurnId, "turn-1"); |
| 124 | first.kill("SIGKILL"); await new Promise((resolve) => first.once("exit", resolve)); |
| 125 | f.items = [{ turn_id: "turn-1", kind: "agent_message", detail: "first second" }]; f.turns[0].status = "completed"; f.seq++; |
| 126 | f.start(); await until(() => f.sent.some((msg) => msg.item_list[0].text_item.text === "first second"), "restart lost final output"); |
| 127 | assert.equal(f.posts.length, 1); assert.equal(f.sent.filter((msg) => msg.item_list[0].text_item.text === "first second").length, 1); |
| 128 | }); |
| 129 | |
| 130 | test("ambiguous admission uses stable operation lookup rather than a duplicate POST", async (t) => { |
| 131 | const f = await fixture(t, { admission: (_req, response) => response.destroy() }); f.batches.push([incoming(1, "work")]); f.start(); |
| 132 | await until(() => f.streams.size, "operation lookup did not recover accepted turn"); |
| 133 | assert.equal(f.posts.length, 1); assert.equal(f.lookups, 1); assert.match(f.posts[0].operation_key, /^[a-f0-9]{64}$/); |
| 134 | assert.equal((await f.disk()).chats.alice.activeTurnId, "turn-1"); |
| 135 | }); |
| 136 | |
| 137 | test("proven rejection retries a stable output ID; uncertain send is retained without resending", async (t) => { |
| 138 | const f = await fixture(t, { send: (response, msg, state, json) => { |
| 139 | if (state.sent.length === 1) return json(response, { ret: -14 }); |
| 140 | return response.destroy(); |
| 141 | } }); f.batches.push([incoming(1, "work")]); f.start(); |
| 142 | await until(() => f.streams.size, "stream missing"); f.event("item.delta", { kind: "agent_message", delta: "retained answer" }); f.complete("retained answer"); |
| 143 | await until(async () => (await f.disk()).chats?.alice?.turnDelivery?.outputs?.[0]?.status === "uncertain", "uncertain send was discarded"); |
| 144 | assert.equal(f.sent.length, 2); assert.equal(f.sent[0].client_id, f.sent[1].client_id); |
| 145 | await delay(1200); assert.equal(f.sent.length, 2); assert.equal((await f.disk()).chats.alice.activeTurnId, "turn-1"); |
| 146 | }); |
| 147 | |
| 148 | test("changed bot account cannot reuse or flush a retained old conversation", async (t) => { |
| 149 | const f = await fixture(t); f.batches.push([incoming(1, "work")]); const first = f.start(); |
| 150 | await until(() => f.streams.size, "stream missing"); first.kill("SIGKILL"); await new Promise((resolve) => first.once("exit", resolve)); |
| 151 | f.items = [{ turn_id: "turn-1", kind: "agent_message", detail: "old private answer" }]; f.turns[0].status = "completed"; |
| 152 | await fs.writeFile(path.join(f.dir, "account.json"), JSON.stringify({ ...f.account, accountId: "bot-B" })); |
| 153 | f.batches.push([incoming(2, "new account prompt")]); f.start(); |
| 154 | await until(() => f.sent.some((msg) => msg.item_list[0].text_item.text.includes("pending turn or retained reply")), "replacement account prompt was not handled"); |
| 155 | assert.equal(f.posts.length, 1); assert.equal(f.sent.some((msg) => msg.item_list[0].text_item.text.includes("old private answer")), false); |
| 156 | assert.equal((await f.disk()).chats.alice.turnDelivery.accountId, "bot-A"); |
| 157 | }); |
| 158 | |
| 159 | test("both Engine capability flags are required before reconciling an uncertain admission", async (t) => { |
| 160 | const f = await fixture(t, { capabilities: { turn_operation_idempotency: true, turn_operation_lookup: false }, admission: (_req, response) => response.destroy() }); |
| 161 | f.batches.push([incoming(1, "legacy admission")]); f.start(); |
| 162 | await until(async () => (await f.disk()).chats?.alice?.pendingAdmission?.submitted, "uncertain request was not retained"); |
| 163 | const polls = f.polls; await until(() => f.polls > polls + 3, "uncertain admission stopped inbound polling"); |
| 164 | await delay(1100); assert.equal(f.posts.length, 1); assert.equal(f.lookups, 0); assert.equal(f.posts[0].operation_key, undefined); |
| 165 | assert.equal((await f.disk()).chats.alice.pendingAdmission.request.prompt, "legacy admission"); |
| 166 | }); |
| 167 | |
| 168 | test("pre-admission inflight claim recovers through the actual production store", async (t) => { |
| 169 | const f = await fixture(t); const key = JSON.stringify(["bot-A", "alice", "1"]); |
| 170 | await fs.writeFile(path.join(f.dir, "thread-map.json"), JSON.stringify({ chats: {}, messages: [key], inflight: { [key]: { at: "fixture" } } })); |
| 171 | f.batches.push([incoming(1, "recover claimed prompt")]); f.start(); |
| 172 | await until(() => f.streams.size, "inflight prompt was silently skipped"); |
| 173 | assert.equal(f.posts.length, 1); assert.equal((await f.disk()).inflight[key], undefined); |
| 174 | assert.equal((await f.disk()).chats.alice.turnDelivery.messageKey, key); |
| 175 | }); |
| 176 | |
| 177 | test("pending new-thread and turn descriptors stay frozen across restart configuration changes", async (t) => { |
| 178 | const f = await fixture(t); const key = JSON.stringify(["bot-A", "alice", "1"]); |
| 179 | const request = { prompt: "frozen prompt", input_summary: "frozen prompt", model: "frozen-model", mode: "agent", allow_shell: false, trust_mode: false, auto_approve: false, operation_key: "frozen-operation" }; |
| 180 | const threadRequest = { model: "frozen-model", workspace: "/frozen/workspace", mode: "agent", allow_shell: false, trust_mode: false, auto_approve: false, archived: false, system_prompt: "frozen system prompt" }; |
| 181 | await fs.writeFile(path.join(f.dir, "thread-map.json"), JSON.stringify({ chats: { alice: { pendingAdmission: { accountId: "bot-A", actorId: "weixin-account:bot-A:user:alice", messageKey: key, request, threadRequest, operationKey: request.operation_key, threadId: null, submitted: false, sinceSeq: null } } }, messages: [], inflight: {} })); |
| 182 | f.start({ CODEWHALE_WORKSPACE: "/changed/workspace", CODEWHALE_MODEL: "changed-model", CODEWHALE_ALLOW_SHELL: "true" }); |
| 183 | await until(() => f.streams.size, "pending request did not recover"); |
| 184 | assert.deepEqual(f.threadRequest, threadRequest); assert.deepEqual(f.posts[0], request); assert.equal(f.posts.length, 1); |
| 185 | }); |
| 186 | |
| 187 | test("completed accepted conversation is not reused by a replacement bot account", async (t) => { |
| 188 | const f = await fixture(t); f.batches.push([incoming(1, "work")]); const first = f.start(); |
| 189 | await until(() => f.streams.size, "stream missing"); f.event("item.delta", { kind: "agent_message", delta: "old answer" }); f.complete("old answer"); |
| 190 | await until(async () => (await f.disk()).chats?.alice?.turnDelivery?.terminal && !(await f.disk()).chats.alice.activeTurnId, "old turn did not settle"); |
| 191 | first.kill("SIGKILL"); await new Promise((resolve) => first.once("exit", resolve)); |
| 192 | await fs.writeFile(path.join(f.dir, "account.json"), JSON.stringify({ ...f.account, accountId: "bot-B" })); |
| 193 | f.batches.push([incoming(2, "replacement prompt")]); f.start(); |
| 194 | await until(() => f.sent.some((msg) => msg.item_list[0].text_item.text.includes("another or unverified bot account")), "completed old binding was silently reused"); |
| 195 | assert.equal(f.posts.length, 1); assert.equal((await f.disk()).chats.alice.bindingAccountId, "bot-A"); |
| 196 | }); |
| 197 | |
| 198 | test("canceled turn retains and sends partial assistant output with terminal status", async (t) => { |
| 199 | const f = await fixture(t); f.batches.push([incoming(1, "work")]); f.start(); |
| 200 | await until(() => f.streams.size, "stream missing"); f.event("item.delta", { kind: "agent_message", delta: "partial answer" }); f.complete("partial answer", "canceled"); |
| 201 | await until(() => f.sent.some((msg) => msg.item_list[0].text_item.text === "partial answer\n\nTurn canceled."), "cancellation dropped partial output"); |
| 202 | assert.equal(f.posts.length, 1); |
| 203 | }); |
| 204 | |
| 205 | test("unanswered admission response leaves polling and status controls live", async (t) => { |
| 206 | const f = await fixture(t, { admission: () => {} }); f.batches.push([incoming(1, "slow admission")]); f.start(); |
| 207 | await until(() => f.posts.length === 1, "admission request missing"); |
| 208 | const started = Date.now(); const polls = f.polls; f.batches.push([incoming(2, "/status")]); |
| 209 | await until(() => f.sent.some((msg) => msg.item_list[0].text_item.text.includes("user_id=alice")), "admission blocked status control until timeout", 3500); |
| 210 | await until(() => f.polls > polls, "inbound polling did not advance while admission remained unanswered", 3500); |
| 211 | assert.ok(Date.now() - started < 7000); assert.equal(f.posts.length, 1); |
| 212 | assert.equal((await f.disk()).chats.alice.pendingAdmission.submitted, true); |
| 213 | }); |
| 214 |