返回 CodeWhale
index.mjs
1 import { WSClient, generateReqId } from "@wecom/aibot-node-sdk";
2
3 import {
4 activeTurnBlock,
5 commandAction,
6 compactRuntimeError,
7 helpText,
8 incomingIdentity,
9 isAllowed,
10 latestRunningTurn,
11 pairingRefusalText,
12 parseBool,
13 parseCommand,
14 parseList,
15 parseApprovalDecisionArgs,
16 isApprovalResponse,
17 isDenyResponse,
18 preservedChatStateFields,
19 requiredEnv,
20 splitMessage,
21 stripGroupPrefix,
22 ThreadStore
23 } from "./lib.mjs";
24 import {
25 ApprovalOwnershipError,
26 createRuntimeClient,
27 decideApproval as decideRuntimeApproval,
28 readJsonSafe,
29 readSse
30 } from "../../bridge-core/src/lib.mjs";
31
32 /** Map of chatId -> latest pending approval info for natural-language approval. */
33 const pendingApprovals = new Map();
34 // Clean up stale approvals every 2 minutes
35 setInterval(() => {
36 const now = Date.now();
37 for (const [chatId, approval] of pendingApprovals) {
38 if (now - approval.timestamp > 300_000) pendingApprovals.delete(chatId);
39 }
40 }, 120_000);
41
42 const config = {
43 botId: requiredEnv("WECOM_BOT_ID"),
44 botSecret: requiredEnv("WECOM_BOT_SECRET"),
45 runtimeUrl: (process.env.CODEWHALE_RUNTIME_URL || "http://127.0.0.1:7878").replace(/\/+$/, ""),
46 runtimeToken: requiredEnv("CODEWHALE_RUNTIME_TOKEN"),
47 workspace: process.env.CODEWHALE_WORKSPACE || process.cwd(),
48 model: process.env.CODEWHALE_MODEL || "auto",
49 mode: process.env.CODEWHALE_MODE || "agent",
50 allowShell: parseBool(process.env.CODEWHALE_ALLOW_SHELL, true),
51 trustMode: parseBool(process.env.CODEWHALE_TRUST_MODE, false),
52 autoApprove: parseBool(process.env.CODEWHALE_AUTO_APPROVE, false),
53 allowlist: parseList(process.env.WECOM_CHAT_ALLOWLIST),
54 allowUnlisted: parseBool(process.env.WECOM_ALLOW_UNLISTED, false),
55 threadMapPath: process.env.WECOM_THREAD_MAP_PATH || "/var/lib/codewhale-wecom-bridge/thread-map.json",
56 maxReplyChars: Number(process.env.WECOM_MAX_REPLY_CHARS || 3500),
57 turnTimeoutMs: Number(process.env.CODEWHALE_TURN_TIMEOUT_MS || 900000),
58 approvalTimeoutMs: Number(process.env.CODEWHALE_APPROVAL_TIMEOUT_MS || 300000)
59 };
60
61 const { runtimeJson, authHeaders } = createRuntimeClient(config);
62
63 const threadStore = await ThreadStore.open(config.threadMapPath);
64
65 const client = new WSClient({
66 botId: config.botId,
67 secret: config.botSecret
68 });
69
70 client.on("message", async (frame) => {
71 try {
72 await handleIncomingMessage(frame);
73 } catch (error) {
74 await reportHandlerError(frame, "Failed to handle WeCom message", error);
75 }
76 });
77
78 client.on("event", async (frame) => {
79 try {
80 await handleEvent(frame);
81 } catch (error) {
82 await reportHandlerError(frame, "Failed to handle WeCom event", error);
83 }
84 });
85
86 client.on("error", (error) => {
87 console.error("WeCom client error:", error);
88 });
89
90 console.log("Starting CodeWhale WeCom bridge");
91 console.log(`Runtime: ${config.runtimeUrl}`);
92 console.log(`Workspace: ${config.workspace}`);
93 if (!config.allowlist.length && !config.allowUnlisted) {
94 console.log("No allowlist configured. Incoming chats will receive their IDs and be refused.");
95 }
96
97 client.connect();
98
99 function replyText(frame, text) {
100 const chunks = splitMessage(text, config.maxReplyChars);
101 return chunks.reduce(
102 (chain, chunk) => chain.then(() => client.replyStream(frame, generateReqId("stream"), chunk, true)),
103 Promise.resolve()
104 );
105 }
106
107 async function reportHandlerError(frame, context, error) {
108 console.error(context, error);
109 try {
110 await replyText(frame, `${context}: ${publicBridgeError(error)}`);
111 } catch (replyError) {
112 console.error("Failed to report WeCom bridge error", replyError);
113 }
114 }
115
116 function publicBridgeError(error) {
117 const message = String(error?.message || error || "unknown error");
118 return message.replaceAll(config.runtimeToken, "<redacted>").slice(0, 500);
119 }
120
121 async function handleIncomingMessage(frame) {
122 const body = frame.body || {};
123 const identity = incomingIdentity(body);
124 console.log(`Incoming message: chatId=${identity.chatId} userId=${identity.userId} chatType=${identity.chatType}`);
125 if (!identity.chatId || !identity.messageId) return;
126
127 if (body.msgtype && body.msgtype !== "text") {
128 await replyText(frame, "目前仅支持文本消息。");
129 return;
130 }
131
132 const textContent = body.text?.content || "";
133 const scoped = stripGroupPrefix(textContent, {
134 chatType: identity.chatType,
135 requirePrefix: identity.chatType === "group",
136 prefix: "/ds"
137 });
138 if (!scoped.accepted) return;
139
140 if (!isAllowed(identity, config.allowlist, config.allowUnlisted)) {
141 await replyText(frame, pairingRefusalText(identity));
142 return;
143 }
144
145 const command = parseCommand(scoped.text);
146 await handleCommand(identity.chatId, command, frame);
147 }
148
149 async function handleEvent(frame) {
150 const body = frame.body || {};
151 const eventType = body.event?.eventtype || "";
152 if (eventType === "enter_chat") {
153 const chatId = body.chatid;
154 if (chatId) {
155 await client.replyWelcome(frame, { msgtype: "text", text: { content: "欢迎使用 CodeWhale!发送 /help 查看可用命令。" } });
156 }
157 }
158 }
159
160 async function handleCommand(chatId, command, frame) {
161 const action = commandAction(command);
162 switch (action.kind) {
163 case "help":
164 await replyText(frame, helpText());
165 return;
166 case "status":
167 await sendStatus(chatId, frame);
168 return;
169 case "threads":
170 await sendThreads(chatId, frame);
171 return;
172 case "new_thread": {
173 const state = await ensureThread(chatId);
174 await replyText(frame, `Created thread ${state.threadId}`);
175 return;
176 }
177 case "resume":
178 await resumeThread(chatId, action.threadId, frame);
179 return;
180 case "interrupt":
181 await interruptActiveTurn(chatId, frame);
182 return;
183 case "compact":
184 await compactThread(chatId, frame);
185 return;
186 case "approval":
187 await decideApproval(chatId, action, frame);
188 return;
189 case "set_model":
190 await setChatModel(chatId, action.modelName, frame);
191 return;
192 case "prompt":
193 // Check if this is a natural-language approval/deny response
194 // This map supplies an approval id for convenience only. Runtime's
195 // pending turn and its durable initiating human still decide authority.
196 if (pendingApprovals.has(chatId)) {
197 const pending = pendingApprovals.get(chatId);
198 if (Date.now() - pending.timestamp < config.approvalTimeoutMs) {
199 if (isApprovalResponse(action.prompt)) {
200 const action2 = { kind: "approval", decision: "allow", approvalId: pending.approvalId };
201 if (await decideApproval(chatId, action2, frame)) pendingApprovals.delete(chatId);
202 return;
203 }
204 if (isDenyResponse(action.prompt)) {
205 const action2 = { kind: "approval", decision: "deny", approvalId: pending.approvalId };
206 if (await decideApproval(chatId, action2, frame)) pendingApprovals.delete(chatId);
207 return;
208 }
209 }
210 }
211 await runPrompt(chatId, action.prompt, frame);
212 return;
213 default:
214 await replyText(frame, helpText());
215 }
216 }
217
218 async function ensureThread(chatId, { forceNew = false } = {}) {
219 const existing = await threadStore.getChat(chatId);
220 if (existing?.threadId && !forceNew) return existing;
221
222 const effectiveModel = existing?.model || config.model;
223
224 const thread = await runtimeJson("/v1/threads", {
225 method: "POST",
226 body: {
227 model: effectiveModel,
228 workspace: config.workspace,
229 mode: config.mode,
230 allow_shell: config.allowShell,
231 trust_mode: config.trustMode,
232 auto_approve: config.autoApprove,
233 archived: false,
234 system_prompt:
235 "You are being controlled from a WeCom (企业微信) phone chat. Keep status updates concise. Ask for tool approvals when needed; do not assume mobile messages imply blanket approval."
236 }
237 });
238
239 const state = {
240 ...preservedChatStateFields(existing),
241 threadId: thread.id,
242 lastSeq: 0,
243 activeTurnId: null,
244 updatedAt: new Date().toISOString()
245 };
246 await threadStore.setChat(chatId, state);
247 return state;
248 }
249
250 function turnActor(identity) {
251 return identity?.userId ? `wecom-user:${identity.userId}` : "";
252 }
253
254 async function runPrompt(chatId, prompt, frame) {
255 if (!prompt.trim()) {
256 await replyText(frame, helpText());
257 return;
258 }
259 const actorId = turnActor(incomingIdentity(frame.body));
260 if (!actorId) throw new ApprovalOwnershipError("turn needs its initiating human");
261 const state = await ensureThread(chatId);
262 const effectiveModel = state?.model || config.model;
263 const detail = await runtimeJson(`/v1/threads/${encodeURIComponent(state.threadId)}`);
264 const activeBlock = activeTurnBlock(detail, state);
265 if (activeBlock) {
266 await threadStore.patchChat(chatId, {
267 activeTurnId: activeBlock.turnId,
268 updatedAt: new Date().toISOString()
269 });
270 await replyText(frame, activeBlock.message);
271 return;
272 }
273 if (state.activeTurnId) {
274 await threadStore.patchChat(chatId, { activeTurnId: null });
275 }
276 const sinceSeq = Number(detail.latest_seq || state.lastSeq || 0);
277
278 const turnResponse = await runtimeJson(
279 `/v1/threads/${encodeURIComponent(state.threadId)}/turns`,
280 {
281 method: "POST",
282 body: {
283 prompt,
284 input_summary: prompt.slice(0, 200),
285 model: effectiveModel,
286 mode: config.mode,
287 allow_shell: config.allowShell,
288 trust_mode: config.trustMode,
289 auto_approve: config.autoApprove
290 }
291 }
292 );
293
294 const turnId = turnResponse.turn?.id;
295 // A turn the Runtime accepted keeps streaming even when its origin cannot be
296 // recorded; its approvals then fail closed and are decided from the TUI.
297 if (turnId && actorId) {
298 await threadStore.recordTurnOrigin(chatId, state.threadId, turnId, actorId);
299 }
300 await threadStore.patchChat(chatId, {
301 activeTurnId: turnId || null,
302 lastSeq: sinceSeq,
303 updatedAt: new Date().toISOString()
304 });
305 await replyText(frame, `Started turn ${turnId || "(unknown)"}`);
306
307 try {
308 await streamTurnEvents(chatId, frame, state.threadId, turnId, sinceSeq);
309 } finally {
310 await threadStore.patchChat(chatId, {
311 activeTurnId: null,
312 updatedAt: new Date().toISOString()
313 });
314 }
315 }
316
317 async function streamTurnEvents(chatId, frame, threadId, turnId, sinceSeq) {
318 const controller = new AbortController();
319 const timeout = setTimeout(() => controller.abort(), config.turnTimeoutMs);
320 const streamId = generateReqId("stream");
321 let responseText = "";
322 let latestSeq = sinceSeq;
323
324 try {
325 const response = await fetch(
326 `${config.runtimeUrl}/v1/threads/${encodeURIComponent(threadId)}/events?since_seq=${sinceSeq}`,
327 {
328 headers: authHeaders(),
329 signal: controller.signal
330 }
331 );
332 if (!response.ok) {
333 const body = await readJsonSafe(response);
334 throw new Error(compactRuntimeError(response.status, body));
335 }
336
337 for await (const event of readSse(response)) {
338 if (!event.data) continue;
339 let record;
340 try {
341 record = JSON.parse(event.data);
342 } catch (error) {
343 console.warn("Skipping malformed runtime SSE event:", publicBridgeError(error));
344 continue;
345 }
346 latestSeq = Math.max(latestSeq, Number(record.seq || 0));
347 await threadStore.patchChat(chatId, { lastSeq: latestSeq });
348
349 if (turnId && record.turn_id && record.turn_id !== turnId) continue;
350
351 if (record.event === "item.delta" && record.payload?.kind === "agent_message") {
352 responseText += record.payload.delta || "";
353 await client.replyStream(frame, streamId, responseText, false);
354 }
355
356 if (record.event === "approval.required") {
357 const approval = record.payload || {};
358 const approvalId = approval.approval_id || approval.id;
359 // Track latest pending approval per chat for natural-language responses
360 if (approvalId) {
361 pendingApprovals.set(chatId, {
362 approvalId,
363 toolName: approval.tool_name || "unknown",
364 description: approval.description || "",
365 timestamp: Date.now()
366 });
367 }
368 await replyText(
369 frame,
370 [
371 "审批请求",
372 `tool=${approval.tool_name || "unknown"}`,
373 `approval_id=${approvalId}`,
374 approval.description || "",
375 "",
376 `回复 /allow ${approvalId}`,
377 `回复 /deny ${approvalId}`,
378 "也可以直接回复「允许」或「拒绝」"
379 ]
380 .filter(Boolean)
381 .join("\n")
382 );
383 }
384
385 if (record.event === "turn.completed") {
386 const turn = record.payload?.turn || {};
387 const status = turn.status || "completed";
388 const errorText = turn.error ? `\n${turn.error}` : "";
389 const fallback = status === "completed" ? "Turn completed." : `Turn ${status}.${errorText}`;
390 await client.replyStream(frame, streamId, responseText.trim() || fallback, true);
391 return;
392 }
393
394 if (record.event === "turn.lifecycle") {
395 const turn = record.payload?.turn || {};
396 const status = turn.status || record.payload?.status;
397 if (["failed", "canceled", "interrupted"].includes(status)) {
398 const errorText = turn.error || record.payload?.error;
399 await client.replyStream(frame, streamId, `Turn ${status}.${errorText ? `\n${errorText}` : ""}`, true);
400 return;
401 }
402 }
403 }
404 } catch (error) {
405 if (error.name === "AbortError") {
406 await replyText(frame, `Turn timed out after ${Math.round(config.turnTimeoutMs / 1000)}s.`);
407 return;
408 }
409 throw error;
410 } finally {
411 clearTimeout(timeout);
412 }
413 }
414
415 async function sendStatus(chatId, frame) {
416 const [health, runtimeInfo, workspace] = await Promise.all([
417 runtimeJson("/health", { auth: false }),
418 runtimeJson("/v1/runtime/info"),
419 runtimeJson("/v1/workspace/status")
420 ]);
421 await replyText(
422 frame,
423 [
424 `runtime=${health.status || "unknown"}`,
425 `version=${runtimeInfo.version || "unknown"}`,
426 `bind=${runtimeInfo.bind_host}:${runtimeInfo.port}`,
427 `auth_required=${runtimeInfo.auth_required}`,
428 `workspace=${workspace.workspace}`,
429 `git_repo=${workspace.git_repo}`,
430 workspace.branch ? `branch=${workspace.branch}` : "",
431 `staged=${workspace.staged} unstaged=${workspace.unstaged} untracked=${workspace.untracked}`
432 ]
433 .filter(Boolean)
434 .join("\n")
435 );
436 }
437
438 async function sendThreads(chatId, frame) {
439 const threads = await runtimeJson("/v1/threads/summary?limit=8&include_archived=true");
440 if (!threads.length) {
441 await replyText(frame, "No runtime threads yet.");
442 return;
443 }
444 await replyText(
445 frame,
446 threads
447 .map((thread) => {
448 const status = thread.latest_turn_status || "none";
449 return `${thread.id} [${status}] ${thread.title || thread.preview || ""}`;
450 })
451 .join("\n")
452 );
453 }
454
455 async function resumeThread(chatId, args, frame) {
456 const threadId = args.trim();
457 if (!threadId) {
458 await replyText(frame, "Usage: /resume <thread_id>");
459 return;
460 }
461 const detail = await runtimeJson(`/v1/threads/${encodeURIComponent(threadId)}`);
462 const existing = await threadStore.getChat(chatId);
463 await threadStore.setChat(chatId, {
464 ...preservedChatStateFields(existing),
465 threadId,
466 lastSeq: Number(detail.latest_seq || 0),
467 activeTurnId: null,
468 updatedAt: new Date().toISOString()
469 });
470 await replyText(frame, `Resumed thread ${threadId}`);
471 }
472
473 async function interruptActiveTurn(chatId, frame) {
474 const state = await threadStore.getChat(chatId);
475 if (!state?.threadId) {
476 await replyText(frame, "No runtime thread recorded for this chat.");
477 return;
478 }
479 const detail = await runtimeJson(`/v1/threads/${encodeURIComponent(state.threadId)}`);
480 const runningTurn = latestRunningTurn(detail);
481 const turnId = state.activeTurnId || runningTurn?.id;
482 if (!turnId) {
483 await replyText(frame, "No active turn recorded for this chat.");
484 return;
485 }
486 await runtimeJson(
487 `/v1/threads/${encodeURIComponent(state.threadId)}/turns/${encodeURIComponent(turnId)}/interrupt`,
488 { method: "POST" }
489 );
490 await threadStore.patchChat(chatId, {
491 activeTurnId: turnId,
492 updatedAt: new Date().toISOString()
493 });
494 await replyText(frame, `Interrupt requested for ${turnId}`);
495 }
496
497 async function compactThread(chatId, frame) {
498 const state = await ensureThread(chatId);
499 const result = await runtimeJson(`/v1/threads/${encodeURIComponent(state.threadId)}/compact`, {
500 method: "POST",
501 body: { reason: "phone bridge request" }
502 });
503 await replyText(frame, `Compaction started: ${result.turn?.id || "unknown turn"}`);
504 }
505
506 async function decideApproval(chatId, action, frame) {
507 const decision = action.decision;
508 // WeCom approvals are text: a session-wide "remember" is never taken from
509 // a chat message.
510 const { approvalId } =
511 action.approvalId != null ? action : parseApprovalDecisionArgs(action.args);
512 const remember = false;
513 if (!approvalId) {
514 await replyText(frame, `Usage: /${decision} <approval_id>`);
515 return false;
516 }
517 try {
518 await decideRuntimeApproval(runtimeJson, { store: threadStore, chatId, actorId: turnActor(incomingIdentity(frame.body)), approvalId, decision, remember });
519 } catch (error) {
520 if (!(error instanceof ApprovalOwnershipError)) throw error;
521 await replyText(frame, `Approval ${approvalId} is not waiting in this chat.`);
522 return false;
523 }
524
525 // Clear activeTurnId so the user can send follow-up messages
526 // immediately instead of being blocked by activeTurnBlock
527 // while the SSE stream processes the turn cancellation.
528 await threadStore.patchChat(chatId, {
529 activeTurnId: null,
530 updatedAt: new Date().toISOString()
531 });
532
533 await replyText(frame, `Approval ${approvalId}: ${decision}${remember ? " and remember" : ""}`);
534 return true;
535 }
536
537 async function setChatModel(chatId, modelName, frame) {
538 if (!modelName || modelName === "default") {
539 await threadStore.patchChat(chatId, {
540 model: null,
541 updatedAt: new Date().toISOString()
542 });
543 await replyText(frame, `Reset per-chat model. Using bridge default: ${config.model}`);
544 return;
545 }
546 await threadStore.patchChat(chatId, {
547 model: modelName,
548 updatedAt: new Date().toISOString()
549 });
550 await replyText(frame, `Per-chat model set to: ${modelName}`);
551 }
552
552 lines Plain Text