返回 CodeWhale
lib.test.mjs
根目录 / integrations / bridge-core / test / lib.test.mjs
1 import test from "node:test";
2 import assert from "node:assert/strict";
3 import { mkdtemp, readdir, readFile, rm } from "node:fs/promises";
4 import { tmpdir } from "node:os";
5 import path from "node:path";
6
7 import {
8 activeTurnBlock,
9 commandAction,
10 createRuntimeClient,
11 envFirst,
12 parseBool,
13 parseCommand,
14 parseEnvText,
15 parseList,
16 parseTextContent,
17 preservedChatStateFields,
18 readJsonSafe,
19 readSse,
20 splitMessage,
21 stripGroupPrefix,
22 ThreadStore,
23 writeFileDurable
24 } from "../src/lib.mjs";
25
26 test("env and primitive parsers handle bridge env conventions", () => {
27 assert.equal(envFirst({ A: "", B: " value " }, "A", "B"), "value");
28 assert.deepEqual(parseList(" a, b ,, "), ["a", "b"]);
29 assert.equal(parseBool("yes"), true);
30 assert.equal(parseBool("0", true), false);
31 assert.deepEqual(parseEnvText("export A='one'\nB=\"two\"\n# nope"), { A: "one", B: "two" });
32 assert.deepEqual(parseEnvText("A='\nB=\"\nEMPTY=\"\""), { A: "'", B: '"', EMPTY: "" });
33 });
34
35 test("parseTextContent supports plain text and JSON text/content wrappers", () => {
36 assert.equal(parseTextContent("hello"), "hello");
37 assert.equal(parseTextContent(JSON.stringify({ text: "hello" })), "hello");
38 assert.equal(parseTextContent(JSON.stringify({ content: "hello" })), "hello");
39 });
40
41 test("stripGroupPrefix supports direct chat types and prefixed group text", () => {
42 assert.deepEqual(
43 stripGroupPrefix("inspect", {
44 chatType: "private",
45 requirePrefix: true,
46 prefix: "/cw",
47 directChatTypes: ["private"]
48 }),
49 { accepted: true, text: "inspect" }
50 );
51 assert.deepEqual(
52 stripGroupPrefix("/cw inspect", {
53 chatType: "group",
54 requirePrefix: true,
55 prefix: "/cw",
56 directChatTypes: ["private"]
57 }),
58 { accepted: true, text: "inspect" }
59 );
60 });
61
62 test("commands map common actions while menu/start stay opt in", () => {
63 assert.deepEqual(parseCommand("/allow@CodeWhaleBot ap_1 remember", { stripBotMention: true }), {
64 name: "allow",
65 args: "ap_1 remember"
66 });
67 assert.deepEqual(parseCommand("/allow@CodeWhaleBot ap_1 remember"), {
68 name: "allow@codewhalebot",
69 args: "ap_1 remember"
70 });
71 assert.deepEqual(commandAction(parseCommand("/status")), { kind: "status" });
72 assert.deepEqual(commandAction(parseCommand("/menu")), { kind: "prompt", prompt: "/menu" });
73 assert.deepEqual(commandAction(parseCommand("/menu"), { allowMenu: true }), { kind: "menu" });
74 assert.deepEqual(commandAction(parseCommand("/start"), { allowStart: true }), { kind: "help" });
75 });
76
77 test("state/message/runtime helpers preserve bridge behavior", () => {
78 assert.deepEqual(
79 preservedChatStateFields({ model: "m", replyToMessageId: "r", ignored: true }, [
80 "model",
81 "replyToMessageId"
82 ]),
83 { model: "m", replyToMessageId: "r" }
84 );
85 assert.deepEqual(splitMessage("a🧪b", 2), ["a🧪", "b"]);
86 assert.deepEqual(splitMessage("alpha beta gamma", 12), ["alpha beta ", "gamma"]);
87 const fenced = splitMessage("```js\nconst first = 1;\nconst second = 2;\n```\nDone", 24);
88 assert.ok(fenced.length > 1);
89 assert.equal(fenced[0].endsWith("\n```"), true);
90 assert.equal(fenced[1].startsWith("```js\n"), true);
91 assert.equal(fenced.at(-1).includes("Done"), true);
92 for (const chunk of fenced) {
93 assert.ok(Array.from(chunk).length <= 24);
94 assert.equal((chunk.match(/```/g) || []).length % 2, 0);
95 }
96 assert.deepEqual(activeTurnBlock({ turns: [{ id: "t1", status: "queued" }] }), {
97 turnId: "t1",
98 message: "Thread already has active turn t1. Wait for it to finish or send /interrupt."
99 });
100 assert.deepEqual(activeTurnBlock({ turns: [{ status: "in_progress" }] }, null), {
101 turnId: "",
102 message: "Thread already has active turn (unknown). Wait for it to finish or send /interrupt."
103 });
104 });
105
106 test("ThreadStore supports chat state, message dedupe, and action tokens", async () => {
107 const dir = await mkdtemp(path.join(tmpdir(), "codewhale-bridge-core-"));
108 try {
109 const statePath = path.join(dir, "thread-map.json");
110 const store = await ThreadStore.open(statePath, {
111 messageLimit: 2,
112 actions: true,
113 actionLimit: 2
114 });
115
116 await store.setChat("chat-a", { threadId: "thread-a" });
117 assert.equal((await store.getChat("chat-a")).threadId, "thread-a");
118
119 assert.equal(await store.recordMessage("m1"), false);
120 assert.equal(await store.recordMessage("m1"), true);
121 assert.equal(await store.recordMessage("m2"), false);
122 assert.equal(await store.recordMessage("m3"), false);
123 assert.deepEqual(store.data.messages, ["m2", "m3"]);
124
125 const token = await store.putAction({ kind: "resume", threadId: "thread-a" }, { chatId: "chat-a" });
126 assert.equal((await store.getAction(token, { chatId: "chat-a" })).kind, "resume");
127 assert.equal((await store.takeAction(token, { chatId: "chat-a" })).threadId, "thread-a");
128 assert.equal(await store.getAction(token, { chatId: "chat-a" }), null);
129
130 const saved = await ThreadStore.open(statePath, { messageLimit: 2, actions: true });
131 assert.equal((await saved.getChat("chat-a")).threadId, "thread-a");
132 assert.deepEqual(saved.data.messages, ["m2", "m3"]);
133 } finally {
134 await rm(dir, { recursive: true, force: true });
135 }
136 });
137
138 test("action tokens remain distinct when clock and legacy randomness repeat", async (t) => {
139 const dir = await mkdtemp(path.join(tmpdir(), "codewhale-action-tokens-"));
140 try {
141 const store = await ThreadStore.open(path.join(dir, "state.json"), { actions: true });
142 t.mock.method(Date, "now", () => 1);
143 t.mock.method(Math, "random", () => 0.5);
144 const first = await store.putAction({ threadId: "thread-a" }, { chatId: "chat-a" });
145 const second = await store.putAction({ threadId: "thread-b" }, { chatId: "chat-a" });
146 assert.notEqual(first, second);
147 assert.match(first, /^[a-f0-9]{32}$/);
148 assert.equal((await store.getAction(first, { chatId: "chat-a" })).threadId, "thread-a");
149 assert.equal((await store.getAction(second, { chatId: "chat-a" })).threadId, "thread-b");
150 } finally {
151 await rm(dir, { recursive: true, force: true });
152 }
153 });
154
155 test("readJsonSafe tolerates empty and non-JSON bodies", async () => {
156 assert.deepEqual(await readJsonSafe({ text: async () => "" }), {});
157 assert.deepEqual(await readJsonSafe({ text: async () => '{"ok":true}' }), { ok: true });
158 assert.equal(await readJsonSafe({ text: async () => "plain text" }), "plain text");
159 });
160
161 test("readSse reassembles events split across chunks and strips CR", async () => {
162 const response = {
163 body: (async function* () {
164 yield Buffer.from('event: item.delta\ndata: {"seq":1}\n\nevent:');
165 yield Buffer.from(' turn.completed\r\ndata: {"seq":2}\n\n');
166 })()
167 };
168 const events = [];
169 for await (const event of readSse(response)) events.push(event);
170 assert.deepEqual(events, [
171 { event: "item.delta", data: '{"seq":1}' },
172 { event: "turn.completed", data: '{"seq":2}' }
173 ]);
174 });
175
176 test("createRuntimeClient sends bearer auth and surfaces runtime errors", async () => {
177 const calls = [];
178 const originalFetch = globalThis.fetch;
179 globalThis.fetch = async (url, options) => {
180 calls.push({ url: String(url), options });
181 if (String(url).endsWith("/fail")) {
182 return {
183 ok: false,
184 status: 503,
185 text: async () => JSON.stringify({ error: { message: "down" } })
186 };
187 }
188 return { ok: true, status: 200, text: async () => JSON.stringify({ ok: true }) };
189 };
190 try {
191 const { runtimeJson, authHeaders } = createRuntimeClient({
192 runtimeUrl: "http://127.0.0.1:7878",
193 runtimeToken: "token-1"
194 });
195 assert.deepEqual(authHeaders(), { authorization: "Bearer token-1" });
196
197 assert.deepEqual(await runtimeJson("/v1/threads", { method: "POST", body: { a: 1 } }), {
198 ok: true
199 });
200 assert.equal(calls[0].url, "http://127.0.0.1:7878/v1/threads");
201 assert.equal(calls[0].options.method, "POST");
202 assert.equal(calls[0].options.headers.authorization, "Bearer token-1");
203 assert.equal(calls[0].options.headers["content-type"], "application/json");
204 assert.equal(calls[0].options.body, JSON.stringify({ a: 1 }));
205
206 await runtimeJson("/health", { auth: false });
207 assert.equal(calls[1].options.method, "GET");
208 assert.deepEqual(calls[1].options.headers, {});
209
210 await assert.rejects(() => runtimeJson("/fail"), /Runtime API request failed \(503\): down/);
211 } finally {
212 globalThis.fetch = originalFetch;
213 }
214 });
215
216 test("ThreadStore batches rapid saves into coalesced durable writes", async () => {
217 const dir = await mkdtemp(path.join(tmpdir(), "codewhale-bridge-core-"));
218 try {
219 const statePath = path.join(dir, "thread-map.json");
220 const store = await ThreadStore.open(statePath);
221 let writes = 0;
222 const originalWrite = store.writeSnapshot.bind(store);
223 store.writeSnapshot = async () => {
224 writes += 1;
225 return originalWrite();
226 };
227
228 await Promise.all(
229 Array.from({ length: 25 }, (_, index) =>
230 store.setChat(`chat-${index}`, { threadId: `thread-${index}` })
231 )
232 );
233 assert.ok(writes <= 2, `expected coalesced writes, saw ${writes}`);
234
235 const saved = await ThreadStore.open(statePath);
236 assert.equal((await saved.getChat("chat-0")).threadId, "thread-0");
237 assert.equal((await saved.getChat("chat-24")).threadId, "thread-24");
238 } finally {
239 await rm(dir, { recursive: true, force: true });
240 }
241 });
242
243 test("ThreadStore persists numeric cursors", async () => {
244 const dir = await mkdtemp(path.join(tmpdir(), "codewhale-bridge-core-"));
245 try {
246 const statePath = path.join(dir, "thread-map.json");
247 const store = await ThreadStore.open(statePath);
248
249 assert.equal(store.getCursor("telegram.update_offset", 7), 7);
250 assert.equal(await store.setCursor("telegram.update_offset", 42), 42);
251 assert.equal(store.getCursor("telegram.update_offset"), 42);
252
253 const saved = await ThreadStore.open(statePath);
254 assert.equal(saved.getCursor("telegram.update_offset"), 42);
255 } finally {
256 await rm(dir, { recursive: true, force: true });
257 }
258 });
259
260 test("writeFileDurable replaces the file through unique temp names and leaves none behind", async () => {
261 const dir = await mkdtemp(path.join(tmpdir(), "codewhale-bridge-core-"));
262 try {
263 const target = path.join(dir, "sync-buf.txt");
264 // Concurrent writers used to share one fixed `.tmp` path and race on it.
265 await Promise.all(Array.from({ length: 8 }, (_, index) => writeFileDurable(target, `cursor-${index}`)));
266 assert.match(await readFile(target, "utf8"), /^cursor-[0-7]$/);
267 assert.deepEqual(await readdir(dir), ["sync-buf.txt"]);
268 } finally {
269 await rm(dir, { recursive: true, force: true });
270 }
271 });
272
273 test("ThreadStore claims survive a restart: finished messages skip, an in-flight one is reported", async () => {
274 const dir = await mkdtemp(path.join(tmpdir(), "codewhale-bridge-core-"));
275 try {
276 const statePath = path.join(dir, "thread-map.json");
277 const store = await ThreadStore.open(statePath, { messageLimit: 10 });
278 assert.equal(await store.claimMessage("u:1"), "new");
279 await store.completeMessage("u:1");
280 assert.equal(await store.claimMessage("u:2"), "new");
281 // The process dies here, before completeMessage("u:2").
282 const restarted = await ThreadStore.open(statePath, { messageLimit: 10 });
283 assert.equal(await restarted.claimMessage("u:1"), "done");
284 assert.equal(await restarted.claimMessage("u:2"), "interrupted");
285 assert.equal(await restarted.claimMessage("u:2"), "done", "reported once, then an ordinary duplicate");
286 assert.equal(await restarted.claimMessage("u:3"), "new");
287 } finally {
288 await rm(dir, { recursive: true, force: true });
289 }
290 });
291
291 lines Plain Text