import * as Lark from "@larksuiteoapi/node-sdk"; import { type BotMessage, type BotReply, type ChannelHandler, type ChannelInstance, type ChannelTransport, createBotMessage, normalizeReply, shouldHandleGroupMessage } from "../../core.js"; import { FeishuStreamingCardSession } from "./streaming-card.js"; export interface FeishuMention { id?: string; name?: string; key?: string; open_id?: string; user_id?: string; union_id?: string; id_container?: { open_id?: string; user_id?: string; union_id?: string; }; } export interface FeishuMessagePayload { message_id?: string; chat_id?: string; chat_type?: string; content?: string; thread_id?: string; } export interface FeishuSenderPayload { sender_id?: { open_id?: string; user_id?: string; union_id?: string; }; id?: string; } export interface FeishuIncomingEvent { event?: FeishuIncomingEvent; message?: FeishuMessagePayload; sender?: FeishuSenderPayload; mentions?: FeishuMention[]; chatType?: string; chatId?: string; senderId?: string; messageId?: string; content?: string; threadId?: string; } export interface FeishuOutgoingMessage extends BotReply { channel: "feishu"; conversationId: string; replyToMessageId: string; } export interface FeishuReactionHandle { messageId: string; reactionId: string; emojiType: string; } export interface FeishuTransport extends ChannelTransport { addReaction?( reaction: Pick ): | Promise | null | void> | Pick | null | void; removeReaction?(reaction: FeishuReactionHandle): Promise | void; } export interface FeishuChannelConfig { appId?: string; appSecret?: string; domain?: string; encryptKey?: string; verificationToken?: string; debug?: boolean; dedupeWindowMs?: number; thinkingReaction?: { enabled?: boolean; emojiType?: string; }; routing?: { groupRequireMention?: boolean; }; transport?: FeishuTransport; } export interface FeishuHandleResult { handled: boolean; reason?: "ignored" | "filtered"; message?: BotMessage; reply?: FeishuOutgoingMessage | null; } export interface FeishuChannel extends ChannelInstance { name: "feishu"; isStarted(): boolean; getSentMessages(): FeishuOutgoingMessage[]; handleEvent(event: FeishuIncomingEvent): Promise; simulateIncomingText(input: { text: string; senderId?: string; conversationId?: string; isDirectMessage?: boolean; mentions?: FeishuMention[]; }): Promise; } const DEFAULT_THINKING_REACTION_EMOJI_TYPE = "OneSecond"; function parseTextContent(content: string | undefined): string { if (!content) return ""; try { const parsed = JSON.parse(content) as { text?: unknown }; if (typeof parsed.text === "string") { return parsed.text; } } catch { return content; } return content; } function unwrapSource(event: FeishuIncomingEvent): FeishuIncomingEvent { return event.event ?? event; } function normalizeMentions(mentions: FeishuMention[] | undefined): string[] { return (mentions ?? []) .map( (mention) => mention.name ?? mention.id ?? mention.key ?? mention.open_id ?? mention.id_container?.open_id ?? mention.user_id ?? mention.id_container?.user_id ) .filter((value): value is string => Boolean(value)); } function normalizeSenderId(source: FeishuIncomingEvent): string { return ( source.sender?.sender_id?.open_id ?? source.sender?.sender_id?.user_id ?? source.sender?.sender_id?.union_id ?? source.sender?.id ?? source.senderId ?? "unknown-feishu-user" ); } function createSdkClient(config: FeishuChannelConfig): Lark.Client { if (!config.appId || !config.appSecret) { throw new Error("Missing Feishu appId/appSecret."); } return new Lark.Client({ appId: config.appId, appSecret: config.appSecret, appType: Lark.AppType.SelfBuild, domain: (config.domain as Lark.Domain | undefined) ?? Lark.Domain.Feishu }); } export function normalizeFeishuEvent(event: FeishuIncomingEvent): BotMessage | null { const source = unwrapSource(event); const message = source.message; if (!message) return null; // This is the main platform boundary: after this point the app should only // reason about BotMessage, not raw Feishu payloads. const chatType = message.chat_type ?? source.chatType ?? "p2p"; return createBotMessage({ channel: "feishu", conversationId: message.chat_id ?? source.chatId ?? "unknown-feishu-chat", senderId: normalizeSenderId(source), messageId: message.message_id ?? source.messageId ?? `feishu-${Date.now()}`, text: parseTextContent(message.content ?? source.content), isDirectMessage: chatType === "p2p", mentions: normalizeMentions(source.mentions), threadId: message.thread_id ?? source.threadId, raw: event }); } function createOutgoingMessage(message: BotMessage, reply: BotReply): FeishuOutgoingMessage { return { channel: "feishu", conversationId: message.conversationId, replyToMessageId: message.messageId, text: reply.text }; } function describeError(error: unknown): string { if (error instanceof Error) { return error.message; } return String(error); } export function createFeishuChannel( config: FeishuChannelConfig = {}, handler: ChannelHandler = {} ): FeishuChannel { const sentMessages: FeishuOutgoingMessage[] = []; const processedMessageIds = new Map(); const inFlightMessageIds = new Set(); let simulatedMessageCounter = 0; const debugEnabled = config.debug ?? (process.env.PI_BOT_DEBUG_FEISHU === "1" || process.env.PI_BOT_DEBUG_FEISHU === "true"); let started = false; let botOpenId: string | undefined; let wsClient: Lark.WSClient | undefined; let sdkClient: Lark.Client | undefined; let channel: FeishuChannel; const dedupeWindowMs = config.dedupeWindowMs ?? 24 * 60 * 60 * 1000; const thinkingReactionEnabled = config.thinkingReaction?.enabled ?? true; const thinkingReactionEmojiType = config.thinkingReaction?.emojiType ?? DEFAULT_THINKING_REACTION_EMOJI_TYPE; function debugLog(message: string, fields?: Record): void { if (!debugEnabled) return; const prefix = `[pi-bot][feishu][${new Date().toISOString()}] ${message}`; if (!fields || Object.keys(fields).length === 0) { console.log(prefix); return; } console.log(prefix, fields); } function pruneProcessedMessageIds(now: number): void { for (const [messageId, timestamp] of processedMessageIds.entries()) { if (now - timestamp > dedupeWindowMs) { processedMessageIds.delete(messageId); } } } /** Detect markdown elements that benefit from card rendering (same heuristic as openclaw). */ function shouldUseCard(text: string): boolean { // Be slightly permissive: during streaming we may only see the opening fence. return /```/.test(text) || /\|.+\|[\r\n]+\|[-:| ]+\|/.test(text); } async function sendReply( message: BotMessage, reply: string | BotReply | null | void ): Promise { const normalizedReply = normalizeReply(reply); if (!normalizedReply) return null; const outgoingMessage = createOutgoingMessage(message, normalizedReply); sentMessages.push(outgoingMessage); // Tests and local simulations inject a transport. Real Feishu mode falls // back to the SDK client below. if (config.transport?.send) { await config.transport.send(outgoingMessage); } else if (sdkClient) { const useCard = shouldUseCard(normalizedReply.text); if (useCard) { await sdkClient.im.message.reply({ path: { message_id: message.messageId }, data: { msg_type: "interactive", content: JSON.stringify({ config: { wide_screen_mode: true }, elements: [ { tag: "markdown", content: normalizedReply.text } ] }) } }); } else { await sdkClient.im.message.reply({ path: { message_id: message.messageId }, data: { msg_type: "text", content: JSON.stringify({ text: normalizedReply.text }) } }); } } return outgoingMessage; } async function sendStreamingReply(message: BotMessage): Promise { if (!handler.onStreamMessage || !sdkClient || !config.appId || !config.appSecret) { return null; } const session = new FeishuStreamingCardSession( sdkClient, { appId: config.appId, appSecret: config.appSecret, domain: config.domain }, debugLog ); let reply: string | BotReply | null | void = null; let streamText = ""; let started = false; let startPromise: Promise | null = null; let shouldStreamCard: boolean | null = null; const ensureStarted = async () => { if (started) return; if (!startPromise) { startPromise = (async () => { try { await session.start(message.conversationId, "chat_id", { replyToMessageId: message.messageId, replyInThread: Boolean(message.threadId) }); started = true; if (streamText) { await session.update(streamText); } } catch (error) { debugLog("failed to start streaming card, fallback to final reply", { error: String(error), messageId: message.messageId }); } })(); } await startPromise; }; try { reply = await handler.onStreamMessage(message, { onMeta(meta) { debugLog("stream meta", { sessionKey: meta.sessionKey, messageId: message.messageId }); }, onDelta(delta) { // `agent.stream()` forwards real append-only text deltas once available. // Concatenate here instead of fuzzy-merging, otherwise short chunks like // "arr" may be dropped just because they already appeared earlier. streamText += delta; if (shouldStreamCard === null && shouldUseCard(streamText)) { shouldStreamCard = true; // Start a streaming card only for markdown-like content. Plain text should // fall back to a normal `text` reply at the end. void ensureStarted(); } if (!started || !session.isActive()) return; void session.update(streamText).catch((error) => { debugLog("failed to update streaming card", { error: String(error), messageId: message.messageId }); }); }, onError(error) { debugLog("stream error", { error, messageId: message.messageId }); } }); const normalizedReply = normalizeReply(reply); const finalText = normalizedReply?.text ?? ""; // Prefer the model's final text to avoid any merge artifacts from streaming. const outputText = finalText || streamText; // Decide card rendering based on the final output if we haven't decided yet. if (shouldStreamCard === null) { shouldStreamCard = shouldUseCard(outputText); } // If we need a card but haven't started yet (e.g. short replies), start now. if (shouldStreamCard && !started) { await ensureStarted(); } if (started && session.isActive()) { await session.close(outputText); if (!normalizedReply) return null; const outgoingMessage = createOutgoingMessage(message, { text: outputText }); sentMessages.push(outgoingMessage); return outgoingMessage; } // No streaming-card started (or not needed): fall back to normal reply path. if (!normalizedReply) { return null; } return await sendReply(message, { text: outputText }); } catch (error) { if (started && session.isActive()) { try { const finalText = normalizeReply(reply)?.text ?? ""; await session.close(finalText || streamText); } catch (closeError) { debugLog("failed to close streaming card after error", { error: String(closeError), messageId: message.messageId }); } } throw error; } } async function addThinkingReaction(message: BotMessage): Promise { if (!thinkingReactionEnabled) return null; try { if (config.transport?.addReaction) { const result = await config.transport.addReaction({ messageId: message.messageId, emojiType: thinkingReactionEmojiType }); if (!result?.reactionId) { debugLog("thinking reaction transport returned no reaction id", { messageId: message.messageId, emojiType: thinkingReactionEmojiType }); return null; } return { messageId: message.messageId, reactionId: result.reactionId, emojiType: result.emojiType ?? thinkingReactionEmojiType }; } if (!sdkClient?.im.messageReaction?.create) { debugLog("thinking reaction api unavailable on sdk client", { messageId: message.messageId, emojiType: thinkingReactionEmojiType }); return null; } const created = await sdkClient.im.messageReaction.create({ path: { message_id: message.messageId }, data: { reaction_type: { emoji_type: thinkingReactionEmojiType } } }); if (!created.data?.reaction_id) { debugLog("thinking reaction create returned no reaction id", { messageId: message.messageId, emojiType: thinkingReactionEmojiType, code: created.code, msg: created.msg }); return null; } debugLog("thinking reaction added", { messageId: message.messageId, reactionId: created.data.reaction_id, emojiType: created.data.reaction_type?.emoji_type ?? thinkingReactionEmojiType }); return { messageId: message.messageId, reactionId: created.data.reaction_id, emojiType: created.data.reaction_type?.emoji_type ?? thinkingReactionEmojiType }; } catch (error) { debugLog("failed to add thinking reaction", { messageId: message.messageId, emojiType: thinkingReactionEmojiType, error: describeError(error) }); return null; } } async function removeThinkingReaction(reaction: FeishuReactionHandle | null): Promise { if (!reaction) return; try { if (config.transport?.removeReaction) { await config.transport.removeReaction(reaction); return; } if (!sdkClient?.im.messageReaction?.delete) { debugLog("thinking reaction delete api unavailable on sdk client", { messageId: reaction.messageId, reactionId: reaction.reactionId, emojiType: reaction.emojiType }); return; } const deleted = await sdkClient.im.messageReaction.delete({ path: { message_id: reaction.messageId, reaction_id: reaction.reactionId } }); debugLog("thinking reaction removed", { messageId: reaction.messageId, reactionId: reaction.reactionId, emojiType: reaction.emojiType, code: deleted.code, msg: deleted.msg }); } catch (error) { debugLog("failed to remove thinking reaction", { messageId: reaction.messageId, reactionId: reaction.reactionId, emojiType: reaction.emojiType, error: describeError(error) }); } } channel = { name: "feishu", async start() { started = true; await config.transport?.start?.(); // If appId/appSecret are present, this channel runs in real Feishu mode. // Feishu pushes events over WS, so no public webhook URL is needed. if (!config.transport && config.appId && config.appSecret) { debugLog("starting real Feishu channel", { domain: config.domain ?? Lark.Domain.Feishu, dedupeWindowMs }); sdkClient = createSdkClient(config); const botInfo = await (sdkClient as unknown as { request(args: { method: string; url: string; data: Record }): Promise<{ bot?: { open_id?: string }; data?: { bot?: { open_id?: string } }; }>; }).request({ method: "GET", url: "/open-apis/bot/v3/info", data: {} }); botOpenId = botInfo?.bot?.open_id ?? botInfo?.data?.bot?.open_id; debugLog("resolved bot identity", { botOpenId }); const dispatcher = new Lark.EventDispatcher({ encryptKey: config.encryptKey ?? "", verificationToken: config.verificationToken ?? "" }); dispatcher.register({ "im.message.receive_v1": async (data: unknown) => { const normalized = normalizeFeishuEvent(data as FeishuIncomingEvent); if (!normalized) return; debugLog("received ws event", { messageId: normalized.messageId, senderId: normalized.senderId, conversationId: normalized.conversationId, isDirectMessage: normalized.isDirectMessage }); // Ignore the bot's own messages to avoid reply loops. if (botOpenId && normalized.senderId === botOpenId) { debugLog("ignored bot self message from dispatcher", { messageId: normalized.messageId, senderId: normalized.senderId }); return; } await channel.handleEvent(data as FeishuIncomingEvent); } } as never); wsClient = new Lark.WSClient({ appId: config.appId, appSecret: config.appSecret, domain: (config.domain as Lark.Domain | undefined) ?? Lark.Domain.Feishu, loggerLevel: Lark.LoggerLevel.warn }); void wsClient.start({ eventDispatcher: dispatcher }); } }, async stop() { started = false; try { wsClient?.close({ force: true }); } catch { // Ignore close errors during shutdown. } wsClient = undefined; await config.transport?.stop?.(); }, isStarted() { return started; }, getSentMessages() { return [...sentMessages]; }, async handleEvent(event) { const message = normalizeFeishuEvent(event); if (!message) { debugLog("ignored event: normalization returned null"); return { handled: false, reason: "ignored" }; } const now = Date.now(); pruneProcessedMessageIds(now); if (botOpenId && message.senderId === botOpenId) { debugLog("ignored bot self message", { messageId: message.messageId, senderId: message.senderId }); return { handled: false, reason: "ignored" }; } if (inFlightMessageIds.has(message.messageId)) { debugLog("ignored duplicate in-flight message", { messageId: message.messageId, senderId: message.senderId, conversationId: message.conversationId }); return { handled: false, reason: "ignored", message }; } if (processedMessageIds.has(message.messageId)) { debugLog("ignored duplicate processed message", { messageId: message.messageId, senderId: message.senderId, conversationId: message.conversationId, firstProcessedAt: processedMessageIds.get(message.messageId) }); return { handled: false, reason: "ignored", message }; } // Group-message filtering happens here so the app runtime only sees // messages that should actually trigger the bot. const shouldHandle = shouldHandleGroupMessage(message, { requireMention: config.routing?.groupRequireMention ?? true }); if (!shouldHandle) { debugLog("filtered group message", { messageId: message.messageId, senderId: message.senderId, conversationId: message.conversationId, mentions: message.mentions }); return { handled: false, reason: "filtered", message }; } inFlightMessageIds.add(message.messageId); debugLog("handling message", { messageId: message.messageId, senderId: message.senderId, conversationId: message.conversationId, text: message.text }); try { const thinkingReaction = await addThinkingReaction(message); let reply: string | BotReply | null | void = null; let outgoingMessage: FeishuOutgoingMessage | null = null; try { if (handler.onStreamMessage && sdkClient && config.appId && config.appSecret) { outgoingMessage = await sendStreamingReply(message); } else { reply = await handler.onMessage?.(message); } } finally { await removeThinkingReaction(thinkingReaction); } if (!outgoingMessage) { outgoingMessage = await sendReply(message, reply); } const processedAt = Date.now(); processedMessageIds.set(message.messageId, processedAt); debugLog("message handled", { messageId: message.messageId, processedAt, replied: Boolean(outgoingMessage), replyLength: outgoingMessage?.text.length ?? 0 }); return { handled: true, message, reply: outgoingMessage }; } finally { inFlightMessageIds.delete(message.messageId); } }, async simulateIncomingText({ text, senderId = "feishu-user", conversationId = "feishu-chat", isDirectMessage = true, mentions = [] }) { simulatedMessageCounter += 1; return this.handleEvent({ message: { message_id: `feishu-sim-${Date.now()}-${simulatedMessageCounter}`, chat_id: conversationId, chat_type: isDirectMessage ? "p2p" : "group", content: JSON.stringify({ text }) }, sender: { sender_id: { open_id: senderId } }, mentions }); } }; return channel; }