返回 CodeWhale
index.mjs
1 import * as Lark from "@larksuiteoapi/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 parseTextContent,
17 preservedChatStateFields,
18 splitMessage,
19 stripGroupPrefix
20 } from "./lib.mjs";
21 import {
22 ApprovalOwnershipError,
23 createRuntimeClient,
24 decideApproval as decideRuntimeApproval,
25 readJsonSafe,
26 readSse,
27 ThreadStore as CoreThreadStore
28 } from "../../bridge-core/src/lib.mjs";
29
30 class ThreadStore extends CoreThreadStore {
31 constructor(filePath) {
32 super(filePath, { messageLimit: 200 });
33 }
34 }
35
36 const config = {
37 appId: requiredEnv("FEISHU_APP_ID"),
38 appSecret: requiredEnv("FEISHU_APP_SECRET"),
39 domain: process.env.FEISHU_DOMAIN || "feishu",
40 runtimeUrl: (process.env.CODEWHALE_RUNTIME_URL || process.env.DEEPSEEK_RUNTIME_URL || "http://127.0.0.1:7878").replace(/\/+$/, ""),
41 runtimeToken: process.env.CODEWHALE_RUNTIME_TOKEN || process.env.DEEPSEEK_RUNTIME_TOKEN || requiredEnv("CODEWHALE_RUNTIME_TOKEN"),
42 workspace: process.env.CODEWHALE_WORKSPACE || process.env.DEEPSEEK_WORKSPACE || process.cwd(),
43 model: process.env.CODEWHALE_MODEL || process.env.DEEPSEEK_MODEL || "auto",
44 mode: process.env.CODEWHALE_MODE || process.env.DEEPSEEK_MODE || "agent",
45 allowShell: parseBool(process.env.CODEWHALE_ALLOW_SHELL ?? process.env.DEEPSEEK_ALLOW_SHELL, true),
46 trustMode: parseBool(process.env.CODEWHALE_TRUST_MODE ?? process.env.DEEPSEEK_TRUST_MODE, false),
47 autoApprove: parseBool(process.env.CODEWHALE_AUTO_APPROVE ?? process.env.DEEPSEEK_AUTO_APPROVE, false),
48 allowlist: parseList(process.env.CODEWHALE_CHAT_ALLOWLIST || process.env.DEEPSEEK_CHAT_ALLOWLIST),
49 allowUnlisted: parseBool(process.env.CODEWHALE_ALLOW_UNLISTED ?? process.env.DEEPSEEK_ALLOW_UNLISTED, false),
50 threadMapPath:
51 process.env.FEISHU_THREAD_MAP_PATH ||
52 "/var/lib/codewhale-feishu-bridge/thread-map.json",
53 allowGroups: parseBool(process.env.FEISHU_ALLOW_GROUPS, false),
54 requirePrefixInGroup: parseBool(process.env.FEISHU_REQUIRE_PREFIX_IN_GROUP, true),
55 groupPrefix: process.env.FEISHU_GROUP_PREFIX || "/ds",
56 maxReplyChars: Number(process.env.FEISHU_MAX_REPLY_CHARS || 3500),
57 turnTimeoutMs: Number(process.env.CODEWHALE_TURN_TIMEOUT_MS || process.env.DEEPSEEK_TURN_TIMEOUT_MS || 900000)
58 };
59
60 const { runtimeJson, authHeaders } = createRuntimeClient(config);
61
62 const sdkConfig = {
63 appId: config.appId,
64 appSecret: config.appSecret,
65 domain: resolveLarkDomain(config.domain)
66 };
67
68 const client = new Lark.Client(sdkConfig);
69 const wsClient = new Lark.WSClient({
70 ...sdkConfig,
71 loggerLevel: Lark.LoggerLevel?.info
72 });
73
74 const threadStore = await ThreadStore.open(config.threadMapPath);
75
76 const dispatcher = new Lark.EventDispatcher({}).register({
77 "im.message.receive_v1": async (data) => {
78 void handleIncomingMessage(data).catch((error) => {
79 console.error("failed to handle incoming Feishu message", error);
80 });
81 }
82 });
83
84 console.log("Starting DeepSeek Feishu bridge");
85 console.log(`Runtime: ${config.runtimeUrl}`);
86 console.log(`Workspace: ${config.workspace}`);
87 if (!config.allowlist.length && !config.allowUnlisted) {
88 console.log("No allowlist configured. Incoming chats will receive their IDs and be refused.");
89 }
90
91 wsClient.start({ eventDispatcher: dispatcher });
92 void reattachActiveTurns().catch((error) => {
93 console.error("failed to reattach active Feishu bridge turns", error);
94 });
95
96 async function handleIncomingMessage(event) {
97 const identity = incomingIdentity(event);
98 if (!identity.chatId) return;
99
100 if (identity.messageType && identity.messageType !== "text") {
101 await sendText(identity.chatId, "Only text messages are supported in this first bridge.");
102 return;
103 }
104
105 const rawText = parseTextContent(event.message?.content || "");
106 const scoped = stripGroupPrefix(rawText, {
107 chatType: identity.chatType,
108 requirePrefix: config.requirePrefixInGroup,
109 prefix: config.groupPrefix
110 });
111 if (!scoped.accepted) return;
112
113 if (identity.messageId && (await threadStore.recordMessage(identity.messageId))) {
114 return;
115 }
116
117 if (identity.chatType !== "p2p" && !config.allowGroups) {
118 await sendText(
119 identity.chatId,
120 "Group chat control is disabled for this bridge. DM the bot, or set FEISHU_ALLOW_GROUPS=true and allowlist this chat."
121 );
122 return;
123 }
124
125 if (!isAllowed(identity, config.allowlist, config.allowUnlisted)) {
126 await sendText(identity.chatId, pairingRefusalText(identity));
127 return;
128 }
129
130 // Only an admitted sender may change delivery/recovery provenance.
131 const { chatId, chatType, openId, unionId, userId } = identity;
132 await threadStore.patchChat(chatId, {
133 authorizedIdentity: { chatId, chatType, openId, unionId, userId },
134 ...(identity.messageId ? { replyToMessageId: identity.messageId } : {}),
135 updatedAt: new Date().toISOString()
136 });
137
138 const command = parseCommand(scoped.text);
139 await handleCommand(identity.chatId, command, identity);
140 }
141
142 async function handleCommand(chatId, command, identity) {
143 const action = commandAction(command);
144 switch (action.kind) {
145 case "help":
146 await sendText(chatId, helpText());
147 return;
148 case "status":
149 await sendStatus(chatId);
150 return;
151 case "threads":
152 await sendThreads(chatId);
153 return;
154 case "new_thread": {
155 const state = await ensureThread(chatId, { forceNew: true });
156 await sendText(chatId, `Created thread ${state.threadId}`);
157 return;
158 }
159 case "resume":
160 await resumeThread(chatId, action.threadId);
161 return;
162 case "interrupt":
163 await interruptActiveTurn(chatId);
164 return;
165 case "compact":
166 await compactThread(chatId);
167 return;
168 case "approval":
169 await decideApproval(chatId, action, identity);
170 return;
171 case "set_model":
172 await setChatModel(chatId, action.modelName);
173 return;
174 case "prompt":
175 await runPrompt(chatId, action.prompt, identity);
176 return;
177 default:
178 await sendText(chatId, helpText());
179 }
180 }
181
182 async function ensureThread(chatId, { forceNew = false } = {}) {
183 const existing = await threadStore.getChat(chatId);
184 if (existing?.threadId && !forceNew) return existing;
185
186 // Use per-chat model if set, fall back to bridge-level default.
187 // / 优先使用 per-chat 模型(/model 命令设置),否则用桥接级别的默认模型。
188 const effectiveModel = existing?.model || config.model;
189
190 const thread = await runtimeJson("/v1/threads", {
191 method: "POST",
192 body: {
193 model: effectiveModel,
194 workspace: config.workspace,
195 mode: config.mode,
196 allow_shell: config.allowShell,
197 trust_mode: config.trustMode,
198 auto_approve: config.autoApprove,
199 archived: false,
200 system_prompt:
201 "You are being controlled from a Feishu/Lark phone chat. Keep status updates concise. Ask for tool approvals when needed; do not assume mobile messages imply blanket approval."
202 }
203 });
204
205 const state = {
206 ...preservedChatStateFields(existing),
207 threadId: thread.id,
208 lastSeq: 0,
209 activeTurnId: null,
210 updatedAt: new Date().toISOString()
211 };
212 await threadStore.setChat(chatId, state);
213 return state;
214 }
215
216 function turnActor(identity) {
217 for (const kind of ["openId", "userId", "unionId"]) {
218 if (identity?.[kind]) return `feishu-${kind}:${identity[kind]}`;
219 }
220 return "";
221 }
222
223 async function runPrompt(chatId, prompt, identity) {
224 if (!prompt.trim()) {
225 await sendText(chatId, helpText());
226 return;
227 }
228 const actorId = turnActor(identity);
229 if (!actorId) throw new ApprovalOwnershipError("turn needs its initiating human");
230 const state = await ensureThread(chatId);
231 // Use per-chat model for this turn (may differ from the thread's
232 // creation model if the user ran /model after the thread was created).
233 // / 使用 per-chat 模型执行本轮对话(如果用户在创建线程后切换过模型)。
234 const effectiveModel = state?.model || config.model;
235 const detail = await runtimeJson(`/v1/threads/${encodeURIComponent(state.threadId)}`);
236 const activeBlock = activeTurnBlock(detail, state);
237 if (activeBlock) {
238 await threadStore.patchChat(chatId, {
239 activeTurnId: activeBlock.turnId,
240 updatedAt: new Date().toISOString()
241 });
242 await sendText(chatId, activeBlock.message);
243 return;
244 }
245 if (state.activeTurnId) {
246 await threadStore.patchChat(chatId, { activeTurnId: null });
247 }
248 const sinceSeq = Number(detail.latest_seq || state.lastSeq || 0);
249
250 const turnResponse = await runtimeJson(
251 `/v1/threads/${encodeURIComponent(state.threadId)}/turns`,
252 {
253 method: "POST",
254 body: {
255 prompt,
256 input_summary: prompt.slice(0, 200),
257 model: effectiveModel,
258 mode: config.mode,
259 allow_shell: config.allowShell,
260 trust_mode: config.trustMode,
261 auto_approve: config.autoApprove
262 }
263 }
264 );
265
266 const turnId = turnResponse.turn?.id;
267 // A turn the Runtime accepted keeps streaming even when its origin cannot be
268 // recorded; its approvals then fail closed and are decided from the TUI.
269 if (turnId && actorId) {
270 await threadStore.recordTurnOrigin(chatId, state.threadId, turnId, actorId);
271 }
272 await threadStore.patchChat(chatId, {
273 activeTurnId: turnId || null,
274 lastSeq: sinceSeq,
275 updatedAt: new Date().toISOString()
276 });
277 await sendText(chatId, `Started turn ${turnId || "(unknown)"}`);
278
279 try {
280 await streamTurnEvents(chatId, state.threadId, turnId, sinceSeq);
281 } finally {
282 await threadStore.patchChat(chatId, {
283 activeTurnId: null,
284 updatedAt: new Date().toISOString()
285 });
286 }
287 }
288
289 async function reattachActiveTurns() {
290 for (const [chatId, state] of threadStore.listChats()) {
291 if (!state?.threadId || !state.activeTurnId) continue;
292 const identity = state.authorizedIdentity;
293 if (!identity || identity.chatId !== chatId ||
294 !["p2p", "group"].includes(identity.chatType) ||
295 (identity.chatType !== "p2p" && !config.allowGroups) ||
296 !isAllowed(identity, config.allowlist, config.allowUnlisted)) continue;
297
298
299 const detail = await runtimeJson(`/v1/threads/${encodeURIComponent(state.threadId)}`);
300 const runningTurn = latestRunningTurn(detail);
301 if (!runningTurn) {
302 await threadStore.patchChat(chatId, {
303 activeTurnId: null,
304 lastSeq: Number(detail.latest_seq || state.lastSeq || 0),
305 updatedAt: new Date().toISOString()
306 });
307 await sendText(chatId, `Bridge restarted. No active turn remains for ${state.threadId}.`);
308 continue;
309 }
310
311 const turnId = runningTurn.id || state.activeTurnId;
312 const sinceSeq = Number(state.lastSeq || 0);
313 await threadStore.patchChat(chatId, {
314 activeTurnId: turnId,
315 updatedAt: new Date().toISOString()
316 });
317 await sendText(
318 chatId,
319 `Bridge restarted. Reattaching to active turn ${turnId} from seq ${sinceSeq}.`
320 );
321 try {
322 await streamTurnEvents(chatId, state.threadId, turnId, sinceSeq);
323 } finally {
324 await threadStore.patchChat(chatId, {
325 activeTurnId: null,
326 updatedAt: new Date().toISOString()
327 });
328 }
329 }
330 }
331
332 async function streamTurnEvents(chatId, threadId, turnId, sinceSeq) {
333 const controller = new AbortController();
334 const timeout = setTimeout(() => controller.abort(), config.turnTimeoutMs);
335 let responseText = "";
336 let latestSeq = sinceSeq;
337 let sentProgressAt = Date.now();
338
339 try {
340 const response = await fetch(
341 `${config.runtimeUrl}/v1/threads/${encodeURIComponent(threadId)}/events?since_seq=${sinceSeq}`,
342 {
343 headers: authHeaders(),
344 signal: controller.signal
345 }
346 );
347 if (!response.ok) {
348 const body = await readJsonSafe(response);
349 throw new Error(compactRuntimeError(response.status, body));
350 }
351
352 for await (const event of readSse(response)) {
353 if (!event.data) continue;
354 const record = JSON.parse(event.data);
355 latestSeq = Math.max(latestSeq, Number(record.seq || 0));
356 await threadStore.patchChat(chatId, { lastSeq: latestSeq });
357
358 if (turnId && record.turn_id && record.turn_id !== turnId) continue;
359
360 if (record.event === "item.delta" && record.payload?.kind === "agent_message") {
361 responseText += record.payload.delta || "";
362 const now = Date.now();
363 if (responseText.length > config.maxReplyChars && now - sentProgressAt > 15000) {
364 await sendText(chatId, responseText.slice(0, config.maxReplyChars));
365 responseText = responseText.slice(config.maxReplyChars);
366 sentProgressAt = now;
367 }
368 }
369
370 if (record.event === "approval.required") {
371 const approval = record.payload || {};
372 await sendText(
373 chatId,
374 [
375 "Approval required",
376 `tool=${approval.tool_name || "unknown"}`,
377 `approval_id=${approval.approval_id || approval.id}`,
378 approval.description || "",
379 "",
380 `Reply /allow ${approval.approval_id || approval.id}`,
381 `Reply /deny ${approval.approval_id || approval.id}`
382 ]
383 .filter(Boolean)
384 .join("\n")
385 );
386 }
387
388 if (record.event === "turn.completed") {
389 const turn = record.payload?.turn || {};
390 const status = turn.status || "completed";
391 const error = turn.error ? `\n${turn.error}` : "";
392 if (status !== "completed") {
393 await sendText(chatId, `Turn ${status}.${error}`.trim());
394 } else {
395 await sendText(chatId, responseText.trim() || "Turn completed.");
396 }
397 return;
398 }
399
400 if (record.event === "turn.lifecycle") {
401 const status = record.payload?.turn?.status || record.payload?.status;
402 if (["failed", "canceled", "interrupted"].includes(status)) {
403 await sendText(chatId, `Turn ${status}.`);
404 return;
405 }
406 }
407 }
408 } catch (error) {
409 if (error.name === "AbortError") {
410 await sendText(chatId, `Turn timed out after ${Math.round(config.turnTimeoutMs / 1000)}s.`);
411 return;
412 }
413 throw error;
414 } finally {
415 clearTimeout(timeout);
416 }
417 }
418
419 async function sendStatus(chatId) {
420 const [health, runtimeInfo, workspace] = await Promise.all([
421 runtimeJson("/health", { auth: false }),
422 runtimeJson("/v1/runtime/info"),
423 runtimeJson("/v1/workspace/status")
424 ]);
425 await sendText(
426 chatId,
427 [
428 `runtime=${health.status || "unknown"}`,
429 `version=${runtimeInfo.version || "unknown"}`,
430 `bind=${runtimeInfo.bind_host}:${runtimeInfo.port}`,
431 `auth_required=${runtimeInfo.auth_required}`,
432 `workspace=${workspace.workspace}`,
433 `git_repo=${workspace.git_repo}`,
434 workspace.branch ? `branch=${workspace.branch}` : "",
435 `staged=${workspace.staged} unstaged=${workspace.unstaged} untracked=${workspace.untracked}`
436 ]
437 .filter(Boolean)
438 .join("\n")
439 );
440 }
441
442 async function sendThreads(chatId) {
443 const threads = await runtimeJson("/v1/threads/summary?limit=8&include_archived=true");
444 if (!threads.length) {
445 await sendText(chatId, "No runtime threads yet.");
446 return;
447 }
448 await sendText(
449 chatId,
450 threads
451 .map((thread) => {
452 const status = thread.latest_turn_status || "none";
453 return `${thread.id} [${status}] ${thread.title || thread.preview || ""}`;
454 })
455 .join("\n")
456 );
457 }
458
459 async function resumeThread(chatId, args) {
460 const threadId = args.trim();
461 if (!threadId) {
462 await sendText(chatId, "Usage: /resume <thread_id>");
463 return;
464 }
465 const detail = await runtimeJson(`/v1/threads/${encodeURIComponent(threadId)}`);
466 const existing = await threadStore.getChat(chatId);
467 await threadStore.setChat(chatId, {
468 ...preservedChatStateFields(existing),
469 threadId,
470 lastSeq: Number(detail.latest_seq || 0),
471 activeTurnId: null,
472 updatedAt: new Date().toISOString()
473 });
474 await sendText(chatId, `Resumed thread ${threadId}`);
475 }
476
477 async function interruptActiveTurn(chatId) {
478 const state = await threadStore.getChat(chatId);
479 if (!state?.threadId) {
480 await sendText(chatId, "No runtime thread recorded for this chat.");
481 return;
482 }
483 const detail = await runtimeJson(`/v1/threads/${encodeURIComponent(state.threadId)}`);
484 const runningTurn = latestRunningTurn(detail);
485 const turnId = state.activeTurnId || runningTurn?.id;
486 if (!turnId) {
487 await sendText(chatId, "No active turn recorded for this chat.");
488 return;
489 }
490 await runtimeJson(
491 `/v1/threads/${encodeURIComponent(state.threadId)}/turns/${encodeURIComponent(
492 turnId
493 )}/interrupt`,
494 { method: "POST" }
495 );
496 await threadStore.patchChat(chatId, {
497 activeTurnId: turnId,
498 updatedAt: new Date().toISOString()
499 });
500 await sendText(chatId, `Interrupt requested for ${turnId}`);
501 }
502
503 async function compactThread(chatId) {
504 const state = await ensureThread(chatId);
505 const result = await runtimeJson(`/v1/threads/${encodeURIComponent(state.threadId)}/compact`, {
506 method: "POST",
507 body: { reason: "phone bridge request" }
508 });
509 await sendText(chatId, `Compaction started: ${result.turn?.id || "unknown turn"}`);
510 }
511
512 async function decideApproval(chatId, action, identity) {
513 const decision = action.decision;
514 const { approvalId, remember } =
515 action.approvalId != null ? action : parseApprovalDecisionArgs(action.args);
516 if (!approvalId) {
517 await sendText(chatId, `Usage: /${decision} <approval_id>${decision === "allow" ? " [remember]" : ""}`);
518 return;
519 }
520 try {
521 await decideRuntimeApproval(runtimeJson, { store: threadStore, chatId, actorId: turnActor(identity), approvalId, decision, remember });
522 } catch (error) {
523 if (!(error instanceof ApprovalOwnershipError)) throw error;
524 await sendText(chatId, `Approval ${approvalId} is not waiting in this chat.`);
525 return;
526 }
527 await sendText(chatId, `Approval ${approvalId}: ${decision}${remember ? " and remember" : ""}`);
528 }
529
530 async function setChatModel(chatId, modelName) {
531 // /model <name> — set per-chat model; "default" or empty resets to bridge default.
532 // / /model "default" 或空参数 — 恢复桥接级别的默认模型。
533 if (!modelName || modelName === "default") {
534 await threadStore.patchChat(chatId, {
535 model: null,
536 updatedAt: new Date().toISOString()
537 });
538 await sendText(chatId, `Reset per-chat model. Using bridge default: ${config.model}`);
539 return;
540 }
541 await threadStore.patchChat(chatId, {
542 model: modelName,
543 updatedAt: new Date().toISOString()
544 });
545 await sendText(chatId, `Per-chat model set to: ${modelName}`);
546 }
547
548 async function sendText(chatId, text) {
549 // Try reply API first — keeps bot responses inside the same Feishu
550 // thread/topic instead of spawning new standalone topics.
551 // / 优先使用 reply API,确保 bot 回复留在话题群的同一条话题内。
552 const state = await threadStore.getChat(chatId);
553 const replyToMessageId = state?.replyToMessageId || null;
554
555 const replyMessage =
556 replyToMessageId
557 ? client.im?.v1?.message?.reply?.bind(client.im.v1.message) ||
558 client.im?.message?.reply?.bind(client.im.message)
559 : null;
560 const createMessage =
561 client.im?.v1?.message?.create?.bind(client.im.v1.message) ||
562 client.im?.message?.create?.bind(client.im.message);
563 if (!createMessage) {
564 throw new Error("Lark SDK client does not expose im message create API");
565 }
566
567 let canReply = Boolean(replyMessage);
568 for (const chunk of splitMessage(text, config.maxReplyChars)) {
569 const body = {
570 msg_type: "text",
571 content: JSON.stringify({ text: chunk })
572 };
573 if (canReply) {
574 try {
575 await replyMessage({
576 path: { message_id: replyToMessageId },
577 data: body
578 });
579 continue;
580 } catch (error) {
581 canReply = false;
582 console.warn("Feishu reply API failed; falling back to message create", error);
583 }
584 }
585 await createMessage({
586 params: { receive_id_type: "chat_id" },
587 data: { ...body, receive_id: chatId }
588 });
589 }
590 }
591
592 function requiredEnv(name) {
593 const value = process.env[name];
594 if (!value || !value.trim()) {
595 throw new Error(`${name} is required`);
596 }
597 return value.trim();
598 }
599
600 function resolveLarkDomain(domain) {
601 const normalized = String(domain || "feishu").toLowerCase();
602 if (normalized === "lark") return Lark.Domain?.Lark || "https://open.larksuite.com";
603 if (normalized === "feishu") return Lark.Domain?.Feishu || "https://open.feishu.cn";
604 return domain;
605 }
606
606 lines Plain Text