返回 CodeWhale
index.mjs
1 import fs from "node:fs/promises";
2 import path from "node:path";
3 import crypto from "node:crypto";
4
5 import {
6 getLoginQR,
7 waitForLogin,
8 getUpdates,
9 sendMessage,
10 getConfig,
11 notifyStart,
12 notifyStop,
13 ILinkLoginBase,
14 parseList,
15 parseBool,
16 envFirst,
17 extractText,
18 parseCommand,
19 commandAction,
20 preservedChatStateFields,
21 splitMessage,
22 compactRuntimeError,
23 latestRunningTurn,
24 helpText,
25 } from "./lib.mjs";
26 import { renderQrToText } from "./qr.mjs";
27 import { ThreadStore as CoreThreadStore, writeFileDurable, ApprovalOwnershipError, decideApproval as decideRuntimeApproval } from "../../bridge-core/src/lib.mjs";
28
29 // ============================================================================
30 // 账号持久化
31 // ============================================================================
32
33 function resolveAccountPath(stateDir) {
34 return path.join(stateDir, "account.json");
35 }
36
37 async function loadAccount(stateDir) {
38 const p = resolveAccountPath(stateDir);
39 try {
40 const raw = await fs.readFile(p, "utf8");
41 return JSON.parse(raw);
42 } catch (error) {
43 if (error.code !== "ENOENT") throw error;
44 return null;
45 }
46 }
47
48 async function saveAccount(stateDir, account) {
49 const p = resolveAccountPath(stateDir);
50 await fs.mkdir(path.dirname(p), { recursive: true, mode: 0o700 });
51 await writeFileDurable(p, `${JSON.stringify(account, null, 2)}\n`, { mode: 0o600 });
52 }
53
54 // ============================================================================
55 // 配置
56 // ============================================================================
57
58 function requiredEnv(name) {
59 const value = process.env[name];
60 if (!value || !value.trim()) {
61 console.error(`Missing required env: ${name}`);
62 process.exit(1);
63 }
64 return value.trim();
65 }
66
67 function requiredEnvFirst(...names) {
68 const value = envFirst(process.env, ...names);
69 if (!value) {
70 console.error(`Missing required env: one of ${names.join(", ")}`);
71 process.exit(1);
72 }
73 return value;
74 }
75
76 // WEIXIN_* is the canonical spelling; the historical WEXIN_* typo names are
77 // still honored as deprecated aliases (each warns once per process).
78 const warnedEnvAliases = new Set();
79
80 function weixinEnv(name) {
81 const legacy = `WEXIN_${name.slice("WEIXIN_".length)}`;
82 const primary = envFirst(process.env, name);
83 if (primary) return primary;
84 const fallback = envFirst(process.env, legacy);
85 if (fallback && !warnedEnvAliases.has(legacy)) {
86 warnedEnvAliases.add(legacy);
87 console.warn(`${legacy} is deprecated; rename it to ${name}.`);
88 }
89 return fallback;
90 }
91
92 const config = {
93 runtimeUrl: (
94 envFirst(process.env, "CODEWHALE_RUNTIME_URL", "DEEPSEEK_RUNTIME_URL") ||
95 "http://127.0.0.1:7878"
96 ).replace(/\/+$/, ""),
97 runtimeToken: requiredEnvFirst(
98 "CODEWHALE_RUNTIME_TOKEN",
99 "DEEPSEEK_RUNTIME_TOKEN"
100 ),
101 workspace:
102 envFirst(process.env, "CODEWHALE_WORKSPACE", "DEEPSEEK_WORKSPACE") ||
103 process.cwd(),
104 model:
105 envFirst(process.env, "CODEWHALE_MODEL", "DEEPSEEK_MODEL") || "auto",
106 mode:
107 envFirst(process.env, "CODEWHALE_MODE", "DEEPSEEK_MODE") || "agent",
108 allowShell: parseBool(
109 envFirst(
110 process.env,
111 "CODEWHALE_ALLOW_SHELL",
112 "DEEPSEEK_ALLOW_SHELL"
113 ),
114 true
115 ),
116 trustMode: parseBool(
117 envFirst(
118 process.env,
119 "CODEWHALE_TRUST_MODE",
120 "DEEPSEEK_TRUST_MODE"
121 ),
122 false
123 ),
124 autoApprove: parseBool(
125 envFirst(
126 process.env,
127 "CODEWHALE_AUTO_APPROVE",
128 "DEEPSEEK_AUTO_APPROVE"
129 ),
130 false
131 ),
132 allowlist: parseList(
133 weixinEnv("WEIXIN_CHAT_ALLOWLIST") ||
134 envFirst(
135 process.env,
136 "CODEWHALE_CHAT_ALLOWLIST",
137 "DEEPSEEK_CHAT_ALLOWLIST"
138 )
139 ),
140 allowUnlisted: parseBool(
141 weixinEnv("WEIXIN_ALLOW_UNLISTED") ||
142 envFirst(
143 process.env,
144 "CODEWHALE_ALLOW_UNLISTED",
145 "DEEPSEEK_ALLOW_UNLISTED"
146 ),
147 false
148 ),
149 stateDir:
150 weixinEnv("WEIXIN_STATE_DIR") ||
151 "/var/lib/codewhale-weixin-bot-bridge",
152 // Defaults inside stateDir rather than to a second absolute path: setting
153 // only WEIXIN_STATE_DIR must not leave the thread map pointing at /var/lib,
154 // which fails with EACCES on every incoming message and silently drops it.
155 threadMapPath:
156 weixinEnv("WEIXIN_THREAD_MAP_PATH") ||
157 path.join(
158 weixinEnv("WEIXIN_STATE_DIR") || "/var/lib/codewhale-weixin-bot-bridge",
159 "thread-map.json"
160 ),
161 maxReplyChars: Number(weixinEnv("WEIXIN_MAX_REPLY_CHARS") || 3500),
162 longPollTimeoutMs: Number(
163 weixinEnv("WEIXIN_LONGPOLL_TIMEOUT_MS") || 35000
164 ),
165 turnTimeoutMs: Number(
166 envFirst(
167 process.env,
168 "CODEWHALE_TURN_TIMEOUT_MS",
169 "DEEPSEEK_TURN_TIMEOUT_MS"
170 ) || 900000
171 ),
172 };
173
174 // ============================================================================
175 // Runtime API 工具
176 // ============================================================================
177
178 function authHeaders() {
179 return {
180 Authorization: `Bearer ${config.runtimeToken}`,
181 "Content-Type": "application/json",
182 };
183 }
184
185 async function readJsonSafe(response) {
186 try {
187 return await response.json();
188 } catch {
189 return null;
190 }
191 }
192
193 async function runtimeJson(subPath, { method = "GET", body = null, auth = true } = {}) {
194 const url = `${config.runtimeUrl}${subPath}`;
195 const options = { method, headers: auth ? authHeaders() : {}, signal: AbortSignal.timeout(10000) };
196 if (body) options.body = JSON.stringify(body);
197 const response = await fetch(url, options);
198 const result = await readJsonSafe(response);
199 if (!response.ok) {
200 const error = new Error(compactRuntimeError(response.status, result));
201 error.status = response.status;
202 throw error;
203 }
204 return result;
205 }
206
207 async function* readSse(response) {
208 let buffer = "";
209 const decoder = new TextDecoder();
210 for await (const chunk of response.body) {
211 buffer += decoder.decode(chunk, { stream: true });
212 const lines = buffer.split("\n");
213 buffer = lines.pop() || "";
214 for (const line of lines) {
215 const trimmed = line.trim();
216 if (!trimmed) continue;
217 if (trimmed.startsWith("data:")) {
218 yield { data: trimmed.slice(5).trim() };
219 } else if (trimmed.startsWith("event:")) {
220 yield { event: trimmed.slice(6).trim() };
221 } else if (trimmed.startsWith("id:")) {
222 yield { id: trimmed.slice(3).trim() };
223 }
224 }
225 }
226 }
227
228 // ============================================================================
229 // 消息发送 — 通过 iLink sendMessage
230 // ============================================================================
231
232 async function sendText(chatId, text, { clientId } = {}) {
233 if (!botAccount) {
234 throw new Error("Bot account unavailable");
235 }
236 const chunks = splitMessage(text, config.maxReplyChars);
237 for (const [index, chunk] of chunks.entries()) {
238 await sendMessage({
239 baseUrl: botAccount.baseUrl,
240 token: botAccount.token,
241 body: {
242 msg: {
243 to_user_id: chatId,
244 client_id: clientId ? `${clientId}-${index}` : crypto.randomUUID(),
245 message_type: 2, // BOT
246 message_state: 2, // FINISH
247 item_list: [{ type: 1, text_item: { text: chunk } }],
248 context_token: await getContextToken(chatId),
249 },
250 },
251 });
252 }
253 }
254
255 async function getContextToken(chatId) {
256 const state = await threadStore.getChat(chatId);
257 return state?.contextToken || undefined;
258 }
259
260 // ============================================================================
261 // 命令处理(与 feishu/telegram/wechat bridge 一致)
262 // ============================================================================
263
264 async function handleCommand(chatId, command, inbound = {}) {
265 const action = commandAction(command);
266 switch (action.kind) {
267 case "help":
268 await sendText(chatId, helpText());
269 return;
270 case "status":
271 await sendStatus(chatId);
272 return;
273 case "threads":
274 await sendThreads(chatId);
275 return;
276 case "new_thread": {
277 const state = await ensureThread(chatId, { forceNew: true });
278 await sendText(chatId, `Created thread ${state.threadId}`);
279 return;
280 }
281 case "resume":
282 await resumeThread(chatId, action.threadId);
283 return;
284 case "interrupt":
285 await interruptActiveTurn(chatId);
286 return;
287 case "compact":
288 await compactThread(chatId);
289 return;
290 case "approval":
291 await decideApproval(chatId, action);
292 return;
293 case "set_model":
294 await setChatModel(chatId, action.modelName);
295 return;
296 case "prompt":
297 await runPrompt(chatId, action.prompt, inbound);
298 return;
299 default:
300 await sendText(chatId, helpText());
301 }
302 }
303
304 async function ensureThread(chatId, { forceNew = false, threadRequest } = {}) {
305 const existing = await threadStore.getChat(chatId);
306 if (forceNew && hasPendingWork(existing)) throw new Error("This chat still has a pending turn or retained reply. Use /status first.");
307 if (existing?.threadId && !forceNew) {
308 if (existing.bindingAccountId !== accountIdentity()) throw new Error("This thread belongs to another or unverified bot account. Use /new to create an explicit binding.");
309 return existing;
310 }
311
312 const effectiveModel = existing?.model || config.model;
313
314 const thread = await runtimeJson("/v1/threads", {
315 method: "POST",
316 body: threadRequest || {
317 model: effectiveModel,
318 workspace: config.workspace,
319 mode: config.mode,
320 allow_shell: config.allowShell,
321 trust_mode: config.trustMode,
322 auto_approve: config.autoApprove,
323 archived: false,
324 system_prompt:
325 "You are being controlled from a WeChat phone chat via iLink Bot. Keep status updates concise. Ask for tool approvals when needed; do not assume mobile messages imply blanket approval.",
326 },
327 });
328
329 const state = {
330 ...preservedChatStateFields(existing, ["model", "authorizedIdentity", "contextToken", "pendingAdmission"]),
331 threadId: thread.id,
332 bindingAccountId: accountIdentity(),
333 lastSeq: 0,
334 activeTurnId: null,
335 updatedAt: new Date().toISOString(),
336 };
337 await threadStore.setChat(chatId, state);
338 return state;
339 }
340
341 function accountIdentity() {
342 return typeof botAccount?.accountId === "string" ? botAccount.accountId.trim() : "";
343 }
344
345 function turnActor(chatId) {
346 return accountIdentity() && chatId ? `weixin-account:${accountIdentity()}:user:${chatId}` : "";
347 }
348
349 function inboundKey(msg) {
350 const id = msg.message_id;
351 if (!accountIdentity() || !msg.from_user_id || id === undefined || id === null || String(id).trim() === "") return null;
352 return JSON.stringify([accountIdentity(), msg.from_user_id, String(id)]);
353 }
354
355 function hasPendingWork(state) {
356 return Boolean(state?.pendingAdmission || state?.activeTurnId || state?.turnDelivery?.outputs?.some((entry) => entry.status !== "accepted"));
357 }
358
359 function ownsRecovery(chatId, state) {
360 const pending = state?.pendingAdmission || state?.turnDelivery;
361 return pending?.accountId === accountIdentity() && pending?.actorId === turnActor(chatId) && isAllowed(chatId);
362 }
363
364 function supportsAdmissionRecovery() {
365 return runtimeCapabilities.turn_operation_idempotency === true && runtimeCapabilities.turn_operation_lookup === true;
366 }
367
368 async function runPrompt(chatId, prompt, { messageKey } = {}) {
369 if (!prompt.trim()) return sendText(chatId, helpText());
370 const current = await threadStore.getChat(chatId);
371 if (hasPendingWork(current)) {
372 startChatRecovery(chatId);
373 return sendText(chatId, "This chat has a pending turn or retained reply. Use /status or /interrupt.");
374 }
375 if (current?.threadId && current.bindingAccountId !== accountIdentity()) return sendText(chatId, "This conversation belongs to another or unverified bot account. Use /new to create a new conversation.");
376 if (!messageKey || !turnActor(chatId)) throw new Error("A stable incoming message identity is required.");
377 const request = {
378 prompt, input_summary: prompt.slice(0, 200), model: current?.model || config.model,
379 mode: config.mode, allow_shell: config.allowShell, trust_mode: config.trustMode, auto_approve: config.autoApprove,
380 };
381 // Retain the frozen request before any admission side effect or poll acknowledgement.
382 const operationKey = supportsAdmissionRecovery() ? crypto.createHash("sha256").update(messageKey).digest("hex") : null;
383 if (operationKey) request.operation_key = operationKey;
384 await threadStore.patchChat(chatId, { pendingAdmission: {
385 accountId: accountIdentity(), actorId: turnActor(chatId), messageKey, request, operationKey,
386 threadId: current?.threadId || null, submitted: false, sinceSeq: null,
387 threadRequest: { model: request.model, workspace: config.workspace, mode: request.mode, allow_shell: request.allow_shell,
388 trust_mode: request.trust_mode, auto_approve: request.auto_approve, archived: false,
389 system_prompt: "You are being controlled from a WeChat phone chat via iLink Bot. Keep status updates concise. Ask for tool approvals when needed; do not assume mobile messages imply blanket approval." },
390 } });
391 startChatRecovery(chatId);
392 }
393
394 async function reconcileAdmission(chatId) {
395 let state = await threadStore.getChat(chatId);
396 let pending = state?.pendingAdmission;
397 if (!pending || !ownsRecovery(chatId, state)) return;
398 if (!pending.threadId) {
399 const thread = await ensureThread(chatId, { threadRequest: pending.threadRequest });
400 // ensureThread replaces chat state only for a new thread; retain the admission.
401 pending = { ...pending, threadId: thread.threadId };
402 await threadStore.patchChat(chatId, { pendingAdmission: pending });
403 }
404 if (pending.sinceSeq === null) {
405 const detail = await runtimeJson(`/v1/threads/${encodeURIComponent(pending.threadId)}`);
406 if (latestRunningTurn(detail)) {
407 await threadStore.patchChat(chatId, { pendingAdmission: { ...pending, blocked: true } });
408 return;
409 }
410 pending = { ...pending, sinceSeq: Number(detail.latest_seq || 0), blocked: false };
411 await threadStore.patchChat(chatId, { pendingAdmission: pending });
412 }
413 let turn;
414 if (pending.submitted) {
415 if (!pending.operationKey || !supportsAdmissionRecovery()) return;
416 try {
417 turn = await runtimeJson(`/v1/threads/${encodeURIComponent(pending.threadId)}/turn-operations/${encodeURIComponent(pending.operationKey)}`);
418 } catch (error) {
419 // 409 is an incomplete durable admission, not permission to create another turn.
420 if (error.status !== 404) throw error;
421 }
422 }
423 if (!turn) {
424 pending = { ...pending, submitted: true };
425 await threadStore.patchChat(chatId, { pendingAdmission: pending });
426 const result = await runtimeJson(`/v1/threads/${encodeURIComponent(pending.threadId)}/turns`, { method: "POST", body: pending.request });
427 turn = result?.turn;
428 }
429 if (!turn?.id || turn.thread_id !== pending.threadId) throw new Error("Invalid Engine admission receipt");
430 await threadStore.recordTurnOrigin(chatId, pending.threadId, turn.id, pending.actorId);
431 await threadStore.patchChat(chatId, {
432 pendingAdmission: null, activeTurnId: turn.id, lastSeq: pending.sinceSeq,
433 turnDelivery: { accountId: pending.accountId, actorId: pending.actorId, threadId: pending.threadId,
434 turnId: turn.id, messageKey: pending.messageKey, nextSeq: pending.sinceSeq, responseText: "", outputs: [], terminal: false },
435 });
436 }
437
438 function queueOutput(delivery, key, text) {
439 if (!text || delivery.outputs.some((entry) => entry.key === key)) return;
440 // Persist one independently acknowledged chunk per API call. A stable client ID
441 // helps reconciliation but iLink does not promise exactly-once recipient delivery.
442 for (const [index, chunk] of splitMessage(text, config.maxReplyChars).entries()) {
443 const chunkKey = `${key}:${index}`;
444 if (delivery.outputs.some((entry) => entry.key === chunkKey)) continue;
445 delivery.outputs.push({ key: chunkKey, text: chunk, status: "queued",
446 clientId: crypto.createHash("sha256").update(`${delivery.accountId}:${delivery.turnId}:${chunkKey}`).digest("hex") });
447 }
448 }
449
450 function finishDelivery(delivery, status) {
451 if (delivery.terminal) return;
452 delivery.terminal = true;
453 delivery.status = status;
454 const text = delivery.responseText.trim();
455 queueOutput(delivery, "result", [text, status !== "completed" ? `Turn ${status}.` : (!text ? "Turn completed." : "")].filter(Boolean).join("\n\n"));
456 }
457
458 async function saveDelivery(chatId, delivery) {
459 await threadStore.patchChat(chatId, { turnDelivery: delivery, lastSeq: delivery.nextSeq });
460 }
461
462 async function flushOutputs(chatId) {
463 let state = await threadStore.getChat(chatId);
464 if (!state?.turnDelivery || !ownsRecovery(chatId, state)) return;
465 let delivery = structuredClone(state.turnDelivery);
466 if (threadStore.turnOrigin(chatId, delivery.threadId, delivery.turnId)?.actorId !== delivery.actorId) return;
467 for (const entry of delivery.outputs) {
468 if (entry.status === "accepted" || entry.status === "uncertain") continue;
469 if (entry.status === "sending") {
470 entry.status = "uncertain";
471 await saveDelivery(chatId, delivery);
472 continue;
473 }
474 entry.status = "sending";
475 await saveDelivery(chatId, delivery);
476 try {
477 await sendText(chatId, entry.text, { clientId: entry.clientId });
478 entry.status = "accepted";
479 } catch (error) {
480 entry.status = error.deliveryStatus === "rejected" ? "queued" : "uncertain";
481 await saveDelivery(chatId, delivery);
482 return;
483 }
484 await saveDelivery(chatId, delivery);
485 }
486 if (delivery.terminal && delivery.outputs.every((entry) => entry.status === "accepted")) {
487 await threadStore.patchChat(chatId, { activeTurnId: null });
488 }
489 }
490
491 function approvalText(approval) {
492 const id = approval.approval_id || approval.id;
493 return ["Approval required", `tool=${approval.tool_name || "unknown"}`, id ? `approval_id=${id}` : "",
494 approval.description || "", id ? `Reply /allow ${id} or /deny ${id}` : "Use the terminal to decide this approval."].filter(Boolean).join("\n");
495 }
496
497 async function reconcileSnapshot(chatId) {
498 const state = await threadStore.getChat(chatId);
499 if (!state?.turnDelivery || !ownsRecovery(chatId, state)) return;
500 const delivery = structuredClone(state.turnDelivery);
501 const detail = await runtimeJson(`/v1/threads/${encodeURIComponent(delivery.threadId)}`);
502 const turn = detail.turns?.find((entry) => entry.id === delivery.turnId && entry.thread_id === delivery.threadId);
503 if (!turn) throw new Error("Accepted Engine turn is unavailable");
504 const items = (detail.items || []).filter((item) => item.turn_id === delivery.turnId && item.kind === "agent_message");
505 if (items.some((item) => typeof item.detail !== "string")) throw new Error("Engine snapshot lacks full assistant output");
506 if (items.length) delivery.responseText = items.map((item) => item.detail).join("\n\n");
507 delivery.nextSeq = Math.max(delivery.nextSeq, Number(detail.latest_seq || 0));
508 for (const approval of detail.pending_approvals || []) {
509 if (approval.turn_id === delivery.turnId) queueOutput(delivery, `approval:${approval.approval_id || approval.id}`, approvalText(approval));
510 }
511 if (["completed", "failed", "canceled", "interrupted"].includes(turn.status)) finishDelivery(delivery, turn.status);
512 await saveDelivery(chatId, delivery);
513 }
514
515 async function streamTurnEvents(chatId, controller) {
516 let state = await threadStore.getChat(chatId);
517 const delivery = structuredClone(state.turnDelivery);
518 // Bound headers and streaming together; this never interrupts the Engine turn.
519 const timeout = setTimeout(() => controller.abort(), config.turnTimeoutMs);
520 try {
521 const response = await fetch(`${config.runtimeUrl}/v1/threads/${encodeURIComponent(delivery.threadId)}/events?since_seq=${delivery.nextSeq}`,
522 { headers: authHeaders(), signal: controller.signal });
523 if (!response.ok) throw new Error("Engine event observation unavailable");
524 for await (const event of readSse(response)) {
525 if (!event.data) continue;
526 const record = JSON.parse(event.data);
527 if (record.event === "stream.end") {
528 if (record.payload?.retryable === false || record.retryable === false) return false;
529 break;
530 }
531 const seq = Number(record.seq || 0);
532 if (seq <= delivery.nextSeq) continue;
533 delivery.nextSeq = seq;
534 if (record.turn_id === delivery.turnId) {
535 if (record.event === "item.delta" && record.payload?.kind === "agent_message") delivery.responseText += record.payload.delta || "";
536 if (record.event === "approval.required") {
537 const approval = record.payload || {};
538 queueOutput(delivery, `approval:${approval.approval_id || approval.id || seq}`, approvalText(approval));
539 }
540 const status = record.payload?.turn?.status || record.payload?.status;
541 if (record.event === "turn.completed") finishDelivery(delivery, status || "completed");
542 if (record.event === "turn.lifecycle" && ["failed", "canceled", "interrupted"].includes(status)) finishDelivery(delivery, status);
543 }
544 // Cursor, accumulator and pending output are committed together before send.
545 await saveDelivery(chatId, delivery);
546 await flushOutputs(chatId);
547 // flushOutputs mutates output states; retain them for the next journal event.
548 delivery.outputs = structuredClone((await threadStore.getChat(chatId)).turnDelivery.outputs);
549 if (delivery.terminal) return true;
550 if (stopping) break;
551 }
552 return true;
553 } finally {
554 clearTimeout(timeout);
555 controller.abort();
556 }
557 }
558
559 function startChatRecovery(chatId) {
560 if (stopping || recoveryTasks.has(chatId)) return;
561 const task = (async () => {
562 while (!stopping) {
563 let state = await threadStore.getChat(chatId);
564 if (!ownsRecovery(chatId, state)) return;
565 try {
566 if (state.pendingAdmission) {
567 await reconcileAdmission(chatId);
568 state = await threadStore.getChat(chatId);
569 if (state.pendingAdmission) {
570 if (state.pendingAdmission.submitted && !state.pendingAdmission.operationKey) return;
571 await sleep(1000);
572 continue;
573 }
574 }
575 if (!state.turnDelivery) return;
576 await reconcileSnapshot(chatId);
577 await flushOutputs(chatId);
578 state = await threadStore.getChat(chatId);
579 if (state.turnDelivery.terminal) {
580 if (!state.turnDelivery.outputs.some((entry) => entry.status === "queued")) return;
581 await sleep(1000);
582 continue;
583 }
584 const controller = new AbortController();
585 recoveryControllers.set(chatId, controller);
586 const retryable = await streamTurnEvents(chatId, controller);
587 recoveryControllers.delete(chatId);
588 if (retryable === false) return;
589 } catch {
590 if (stopping) return;
591 console.warn("Turn observation interrupted; retaining the accepted turn for recovery.");
592 }
593 await sleep(1000);
594 }
595 })().finally(() => { recoveryTasks.delete(chatId); recoveryControllers.delete(chatId); });
596 recoveryTasks.set(chatId, task);
597 }
598
599 async function sendStatus(chatId) {
600 try {
601 const state = await threadStore.getChat(chatId);
602 const [health, runtimeInfo, workspace] = await Promise.all([
603 runtimeJson("/health", { auth: false }),
604 runtimeJson("/v1/runtime/info"),
605 runtimeJson("/v1/workspace/status"),
606 ]);
607 await sendText(
608 chatId,
609 [
610 `user_id=${chatId}`,
611 state?.pendingAdmission ? "Turn admission is pending; its original request is retained." : "",
612 state?.activeTurnId ? `accepted_turn=${state.activeTurnId}` : "",
613 state?.turnDelivery?.outputs?.some((entry) => entry.status === "uncertain") ? "A reply has unconfirmed API acceptance. It is retained for review and will not be resent automatically." : "",
614 `runtime=${health.status || "unknown"}`,
615 `version=${runtimeInfo.version || "unknown"}`,
616 `bind=${runtimeInfo.bind_host}:${runtimeInfo.port}`,
617 `auth_required=${runtimeInfo.auth_required}`,
618 `workspace=${workspace.workspace}`,
619 `git_repo=${workspace.git_repo}`,
620 workspace.branch ? `branch=${workspace.branch}` : "",
621 `staged=${workspace.staged} unstaged=${workspace.unstaged} untracked=${workspace.untracked}`,
622 ]
623 .filter(Boolean)
624 .join("\n")
625 );
626 } catch (error) {
627 await sendText(chatId, `Status check failed: ${error.message}`);
628 }
629 }
630
631 async function sendThreads(chatId) {
632 try {
633 const threads = await runtimeJson(
634 "/v1/threads/summary?limit=8&include_archived=true"
635 );
636 if (!threads.length) {
637 await sendText(chatId, "No runtime threads yet.");
638 return;
639 }
640 await sendText(
641 chatId,
642 threads
643 .map((thread) => {
644 const status = thread.latest_turn_status || "none";
645 return `${thread.id} [${status}] ${thread.title || thread.preview || ""}`;
646 })
647 .join("\n")
648 );
649 } catch (error) {
650 await sendText(chatId, `Thread listing failed: ${error.message}`);
651 }
652 }
653
654 async function resumeThread(chatId, args) {
655 const threadId = args.trim();
656 if (!threadId) {
657 await sendText(chatId, "Usage: /resume <thread_id>");
658 return;
659 }
660 try {
661 const existing = await threadStore.getChat(chatId);
662 if (hasPendingWork(existing)) throw new Error("This chat still has a pending turn or retained reply. Use /status first.");
663 const detail = await runtimeJson(
664 `/v1/threads/${encodeURIComponent(threadId)}`
665 );
666 await threadStore.setChat(chatId, {
667 ...preservedChatStateFields(existing, ["model", "authorizedIdentity", "contextToken"]),
668 threadId,
669 bindingAccountId: accountIdentity(),
670 lastSeq: Number(detail.latest_seq || 0),
671 activeTurnId: null,
672 updatedAt: new Date().toISOString(),
673 });
674 await sendText(chatId, `Resumed thread ${threadId}`);
675 } catch (error) {
676 await sendText(chatId, `Resume failed: ${error.message}`);
677 }
678 }
679
680 async function interruptActiveTurn(chatId) {
681 const state = await threadStore.getChat(chatId);
682 if (!state?.threadId) {
683 await sendText(chatId, "No runtime thread recorded for this chat.");
684 return;
685 }
686 try {
687 const detail = await runtimeJson(
688 `/v1/threads/${encodeURIComponent(state.threadId)}`
689 );
690 const turnId = state.activeTurnId;
691 if (state.turnDelivery?.accountId !== accountIdentity() || threadStore.turnOrigin(chatId, state.threadId, turnId)?.actorId !== turnActor(chatId)) throw new ApprovalOwnershipError("Turn is not owned by this account and chat");
692 if (!turnId) {
693 await sendText(chatId, "No active turn recorded for this chat.");
694 return;
695 }
696 await runtimeJson(
697 `/v1/threads/${encodeURIComponent(state.threadId)}/turns/${encodeURIComponent(turnId)}/interrupt`,
698 { method: "POST" }
699 );
700 await threadStore.patchChat(chatId, {
701 activeTurnId: turnId,
702 updatedAt: new Date().toISOString(),
703 });
704 await sendText(chatId, `Interrupt requested for ${turnId}`);
705 } catch (error) {
706 await sendText(chatId, `Interrupt failed: ${error.message}`);
707 }
708 }
709
710 async function compactThread(chatId) {
711 try {
712 const state = await ensureThread(chatId);
713 const result = await runtimeJson(
714 `/v1/threads/${encodeURIComponent(state.threadId)}/compact`,
715 {
716 method: "POST",
717 body: { reason: "weixin-bot bridge request" },
718 }
719 );
720 await sendText(
721 chatId,
722 `Compaction started: ${result.turn?.id || "unknown turn"}`
723 );
724 } catch (error) {
725 await sendText(chatId, `Compact failed: ${error.message}`);
726 }
727 }
728
729 async function decideApproval(chatId, action) {
730 const decision = action.decision;
731 const { approvalId, remember } = action;
732 if (!approvalId) {
733 await sendText(
734 chatId,
735 `Usage: /${decision} <approval_id>${decision === "allow" ? " [remember]" : ""}`
736 );
737 return;
738 }
739 try {
740 await decideRuntimeApproval(runtimeJson, { store: threadStore, chatId, actorId: turnActor(chatId), approvalId, decision, remember });
741 await sendText(
742 chatId,
743 `Approval ${approvalId}: ${decision}${remember ? " and remember" : ""}`
744 );
745 } catch (error) {
746 if (error instanceof ApprovalOwnershipError) {
747 await sendText(chatId, `Approval ${approvalId} is not waiting in this chat.`);
748 return;
749 }
750 await sendText(chatId, `Approval failed: ${error.message}`);
751 }
752 }
753
754 async function setChatModel(chatId, modelName) {
755 if (!modelName || modelName === "default") {
756 await threadStore.patchChat(chatId, {
757 model: null,
758 updatedAt: new Date().toISOString(),
759 });
760 await sendText(
761 chatId,
762 `Reset per-chat model. Using bridge default: ${config.model}`
763 );
764 return;
765 }
766 await threadStore.patchChat(chatId, {
767 model: modelName,
768 updatedAt: new Date().toISOString(),
769 });
770 await sendText(chatId, `Per-chat model set to: ${modelName}`);
771 }
772
773 // ============================================================================
774 // 主循环 — 长轮询 getUpdates
775 // ============================================================================
776
777 let botAccount = null;
778 let stopping = false;
779 let threadStore;
780 const recoveryTasks = new Map();
781 const recoveryControllers = new Map();
782 let runtimeCapabilities = {};
783
784 function resolveSyncBufPath(stateDir) {
785 return path.join(stateDir, "sync-buf.txt");
786 }
787
788 async function loadSyncBuf(stateDir) {
789 const p = resolveSyncBufPath(stateDir);
790 try {
791 return await fs.readFile(p, "utf8");
792 } catch {
793 return "";
794 }
795 }
796
797 async function saveSyncBuf(stateDir, buf) {
798 const p = resolveSyncBufPath(stateDir);
799 // The state dir may not exist yet on a first run whose first persisted write
800 // is the poll cursor rather than account.json.
801 await fs.mkdir(path.dirname(p), { recursive: true, mode: 0o700 });
802 await writeFileDurable(p, buf, { mode: 0o600 });
803 }
804
805 /** One durably claimed inbound message; agent observation runs separately. */
806 async function handleInbound(msg) {
807 const fromUser = msg.from_user_id || "";
808 if (msg.message_type !== undefined && msg.message_type !== 1) return;
809 if (!isAllowed(fromUser)) return;
810 if (msg.context_token) await threadStore.patchChat(fromUser, { contextToken: msg.context_token });
811 const text = extractText(msg.item_list);
812 if (!text) return sendText(fromUser, "仅支持文本消息。图片/语音/视频/文件暂不支持。");
813 await handleCommand(fromUser, parseCommand(text), { messageKey: inboundKey(msg) });
814 }
815
816 async function monitorLoop() {
817 const { baseUrl, token } = botAccount;
818 let getUpdatesBuf = await loadSyncBuf(config.stateDir);
819 let nextTimeoutMs = config.longPollTimeoutMs;
820 let consecutiveFailures = 0;
821
822 console.log(`Monitor started: baseUrl=${baseUrl} timeoutMs=${nextTimeoutMs}`);
823
824 while (!stopping) {
825 try {
826 const abortController = new AbortController();
827 const timer = setTimeout(
828 () => abortController.abort(),
829 nextTimeoutMs + 5000
830 );
831
832 const resp = await getUpdates({
833 baseUrl,
834 token,
835 get_updates_buf: getUpdatesBuf,
836 timeoutMs: nextTimeoutMs,
837 signal: abortController.signal,
838 });
839
840 clearTimeout(timer);
841
842 if (resp.longpolling_timeout_ms) {
843 nextTimeoutMs = resp.longpolling_timeout_ms;
844 }
845
846 // 检查错误
847 const isApiError =
848 (resp.ret !== undefined && resp.ret !== 0) ||
849 (resp.errcode !== undefined && resp.errcode !== 0);
850
851 if (isApiError) {
852 consecutiveFailures += 1;
853 console.error(
854 `getUpdates error: ret=${resp.ret} errcode=${resp.errcode} errmsg=${resp.errmsg}`
855 );
856 if (consecutiveFailures >= 3) {
857 console.error("3 consecutive failures, backing off 30s");
858 await sleep(30000);
859 consecutiveFailures = 0;
860 } else {
861 await sleep(2000);
862 }
863 continue;
864 }
865
866 consecutiveFailures = 0;
867
868 // Claims stay inflight until durable admission/handling succeeds. A crash
869 // replays commands conservatively and reconciles the persisted prompt receipt.
870 for (const msg of resp.msgs || []) {
871 const key = inboundKey(msg);
872 if (!key) { console.warn("Ignored inbound message without a stable account/peer/message identity."); continue; }
873 // Core claimMessage retires an interrupted claim before returning it.
874 // Keep this caller's inflight marker until its recovery payload is durable.
875 const claim = threadStore.data.inflight?.[key] ? "interrupted" : await threadStore.claimMessage(key);
876 if (claim === "done") continue;
877 if (claim === "interrupted") {
878 const state = await threadStore.getChat(msg.from_user_id);
879 if (state?.pendingAdmission?.messageKey === key || state?.turnDelivery?.messageKey === key) {
880 startChatRecovery(msg.from_user_id);
881 } else if (commandAction(parseCommand(extractText(msg.item_list))).kind === "prompt") {
882 // Prompt admission always persists before side effects. No receipt
883 // means this claim crashed before admission, so it can be handled.
884 await handleInbound(msg);
885 } else if (isAllowed(msg.from_user_id)) {
886 await sendText(msg.from_user_id, "The bridge restarted while handling this command. Its effect is uncertain; check /status before repeating it.");
887 }
888 await threadStore.completeMessage(key);
889 continue;
890 }
891 await handleInbound(msg);
892 await threadStore.completeMessage(key);
893 }
894 if (resp.get_updates_buf) {
895 await saveSyncBuf(config.stateDir, resp.get_updates_buf);
896 getUpdatesBuf = resp.get_updates_buf;
897 }
898 for (const [chatId, state] of threadStore.listChats()) {
899 if (hasPendingWork(state)) startChatRecovery(chatId);
900 }
901 } catch (error) {
902 if (error.name === "AbortError" || error.message?.includes("abort")) {
903 // 长轮询超时是正常的,立即重试
904 continue;
905 }
906 if (stopping) break;
907
908 consecutiveFailures += 1;
909 console.error(
910 `getUpdates exception (${consecutiveFailures}/3):`,
911 error.message
912 );
913 if (consecutiveFailures >= 3) {
914 console.error("3 consecutive exceptions, backing off 30s");
915 await sleep(30000);
916 consecutiveFailures = 0;
917 } else {
918 await sleep(2000);
919 }
920 }
921 }
922 }
923
924 function isAllowed(fromUser) {
925 if (config.allowUnlisted) return true;
926 const allowed = new Set(config.allowlist);
927 return allowed.has(fromUser);
928 }
929
930 function sleep(ms) {
931 return new Promise((resolve) => setTimeout(resolve, ms));
932 }
933
934 // ============================================================================
935 // 启动流程 — QR 登录 → 长轮询
936 // ============================================================================
937
938 async function main() {
939 console.log("Starting CodeWhale Weixin Bot Bridge");
940 console.log(`Runtime: ${config.runtimeUrl}`);
941 console.log(`Workspace: ${config.workspace}`);
942 console.log(`State dir: ${config.stateDir}`);
943 console.log(`Thread map: ${config.threadMapPath}`);
944
945 // 初始化 ThreadStore。`open()` 只读,真正的写入发生在第一条消息到达时;
946 // 那时失败会被 getUpdates 的 catch 吞掉,表现为“微信没有回应”。所以这里
947 // 先建目录并真实写一次探针文件,把问题在启动时就暴露出来。
948 try {
949 const dir = path.dirname(config.threadMapPath);
950 await fs.mkdir(dir, { recursive: true, mode: 0o700 });
951 const probe = path.join(dir, ".write-probe");
952 await fs.writeFile(probe, "", { mode: 0o600 });
953 await fs.rm(probe, { force: true });
954 } catch (error) {
955 console.error(
956 `Thread map directory is not writable: ${path.dirname(config.threadMapPath)} (${error.message})`
957 );
958 console.error(
959 "Set WEIXIN_STATE_DIR (or WEIXIN_THREAD_MAP_PATH) to a writable directory."
960 );
961 process.exit(1);
962 }
963 threadStore = await CoreThreadStore.open(config.threadMapPath, { messageLimit: 500, privateMode: true });
964
965 // 尝试加载已有账号
966 botAccount = await loadAccount(config.stateDir);
967
968 if (botAccount?.token) {
969 console.log("Loaded existing bot account, trying to resume...");
970 console.log(` accountId: ${botAccount.accountId}`);
971 console.log(` baseUrl: ${botAccount.baseUrl}`);
972 } else {
973 // QR 登录
974 console.log("No bot account found. Starting QR login...");
975 console.log("");
976
977 const { qrcodeUrl, sessionKey } = await getLoginQR();
978 console.log("请用微信扫描以下二维码登录:");
979 // Render the login URL as a scannable terminal QR. The URL is printed too,
980 // so a terminal that mangles the half-block glyphs still has a way through.
981 try {
982 if (qrcodeUrl) console.log(renderQrToText(qrcodeUrl));
983 } catch (error) {
984 console.warn(`Could not render QR in terminal: ${error.message}`);
985 }
986 console.log(qrcodeUrl);
987 console.log("");
988
989 const result = await waitForLogin({ sessionKey, timeoutMs: 300_000 });
990
991 if (!result.connected) {
992 console.error(`Login failed: ${result.message}`);
993 process.exit(1);
994 }
995
996 botAccount = {
997 accountId: result.accountId,
998 token: result.botToken,
999 baseUrl: result.baseUrl,
1000 userId: result.userId,
1001 };
1002
1003 await saveAccount(config.stateDir, botAccount);
1004 console.log(`✅ Login successful! accountId=${botAccount.accountId}`);
1005 }
1006
1007 if (!accountIdentity()) throw new Error("The saved bot account has no stable account identity; pair again before receiving messages.");
1008 try { runtimeCapabilities = (await runtimeJson("/v1/runtime/info"))?.capabilities || {}; } catch { console.warn("Engine recovery capabilities unavailable; uncertain admissions will remain retained."); }
1009 for (const [chatId, state] of threadStore.listChats()) {
1010 if (hasPendingWork(state)) startChatRecovery(chatId);
1011 }
1012
1013 // 通知上线
1014 try {
1015 const startResp = await notifyStart({
1016 baseUrl: botAccount.baseUrl,
1017 token: botAccount.token,
1018 });
1019 if (startResp.ret && startResp.ret !== 0) {
1020 console.warn(`notifyStart: ret=${startResp.ret} errmsg=${startResp.errmsg}`);
1021 } else {
1022 console.log("notifyStart: OK");
1023 }
1024 } catch (error) {
1025 console.error("notifyStart failed:", error.message);
1026 }
1027
1028 // 信号处理
1029 process.once("SIGINT", shutdown);
1030 process.once("SIGTERM", shutdown);
1031
1032 if (!config.allowlist.length && !config.allowUnlisted) {
1033 console.log(
1034 "No allowlist configured. Incoming chats will receive their user IDs and be refused."
1035 );
1036 }
1037
1038 // 进入长轮询循环
1039 await monitorLoop();
1040
1041 console.log("Bridge stopped.");
1042 }
1043
1044 async function shutdown() {
1045 if (stopping) return;
1046 stopping = true;
1047 console.log("Shutting down...");
1048 for (const controller of recoveryControllers.values()) controller.abort();
1049
1050 if (botAccount?.token) {
1051 try {
1052 const stopResp = await notifyStop({
1053 baseUrl: botAccount.baseUrl,
1054 token: botAccount.token,
1055 });
1056 console.log(
1057 `notifyStop: ret=${stopResp.ret} errmsg=${stopResp.errmsg ?? "OK"}`
1058 );
1059 } catch (error) {
1060 console.error("notifyStop failed:", error.message);
1061 }
1062 }
1063
1064 setTimeout(() => process.exit(0), 2000);
1065 }
1066
1067 main().catch((error) => {
1068 console.error("Fatal error:", error);
1069 process.exit(1);
1070 });
1071
1071 lines Plain Text