返回 DeepSeek-Reasonix
feishu.go
根目录 / internal / bot / feishu / feishu.go
1 // Package feishu 实现飞书自建应用 Bot 适配器。
2 // 参考 Hermes Agent 的 feishu adapter:
3 // - 长连接 WebSocket(默认)或 Webhook 模式
4 // - @mention gating
5 // - open_id / user_id / union_id 映射
6 // - 消息去重
7 // - interactive card 审批/问答
8 package feishu
9
10 import (
11 "context"
12 "crypto/sha256"
13 "crypto/subtle"
14 "encoding/hex"
15 "encoding/json"
16 "errors"
17 "fmt"
18 "io"
19 "log/slog"
20 "maps"
21 "net/http"
22 "os"
23 "strings"
24 "sync"
25 "time"
26
27 "reasonix/internal/bot"
28 "reasonix/internal/config"
29
30 lark "github.com/larksuite/oapi-sdk-go/v3"
31 larkcore "github.com/larksuite/oapi-sdk-go/v3/core"
32 "github.com/larksuite/oapi-sdk-go/v3/event/dispatcher"
33 "github.com/larksuite/oapi-sdk-go/v3/event/dispatcher/callback"
34 larkcontact "github.com/larksuite/oapi-sdk-go/v3/service/contact/v3"
35 larkim "github.com/larksuite/oapi-sdk-go/v3/service/im/v1"
36 larkws "github.com/larksuite/oapi-sdk-go/v3/ws"
37 )
38
39 // textContent 飞书消息文本内容结构。
40 type textContent struct {
41 Text string `json:"text"`
42 }
43
44 const feishuPendingReactionEmoji = "OnIt"
45
46 // feishuEvent 飞书事件结构。
47 type feishuEvent struct {
48 Schema string `json:"schema"`
49 Header feishuHeader `json:"header"`
50 Event json.RawMessage `json:"event"`
51 }
52
53 type feishuHeader struct {
54 EventID string `json:"event_id"`
55 EventType string `json:"event_type"`
56 Token string `json:"token"`
57 CreateTime string `json:"create_time"`
58 }
59
60 type feishuMsgEvent struct {
61 MessageID string `json:"message_id"`
62 RootID string `json:"root_id"`
63 ParentID string `json:"parent_id"`
64 ThreadID string `json:"thread_id"`
65 ChatID string `json:"chat_id"`
66 ChatType string `json:"chat_type"`
67 MsgType string `json:"msg_type"`
68 Content string `json:"content"`
69 Sender feishuSender `json:"sender"`
70 Mentions []feishuMention `json:"mentions"`
71 }
72
73 type feishuSender struct {
74 SenderID struct {
75 UserID string `json:"user_id"`
76 OpenID string `json:"open_id"`
77 UnionID string `json:"union_id"`
78 } `json:"sender_id"`
79 }
80
81 type feishuMention struct {
82 Key string `json:"key"`
83 Name string `json:"name"`
84 ID struct {
85 OpenID string `json:"open_id"`
86 } `json:"id"`
87 }
88
89 func webhookMentionRefs(mentions []feishuMention) []mentionRef {
90 refs := make([]mentionRef, 0, len(mentions))
91 for _, m := range mentions {
92 refs = append(refs, mentionRef{Key: m.Key, OpenID: m.ID.OpenID, Name: m.Name})
93 }
94 return refs
95 }
96
97 // adapter 飞书适配器实现。
98 type adapter struct {
99 cfg config.FeishuBotConfig
100 logger *slog.Logger
101 msgCh chan bot.InboundMessage
102 cancel context.CancelFunc
103 client *lark.Client
104 wsClient *larkws.Client
105
106 // fetchResource 覆盖消息资源下载(测试注入);nil 时用 sdkFetchResource。
107 fetchResource func(ctx context.Context, messageID, key, typ string) ([]byte, string, error)
108
109 clientMu sync.Mutex // 保护 client 懒初始化
110
111 seenMu sync.Mutex
112 seen map[string]bool // 消息去重
113
114 botMu sync.Mutex
115 botID string // bot 自身 open_id,用于群聊 @ 门控与占位符剔除
116
117 nameMu sync.Mutex
118 names map[string]nameCacheEntry // open_id -> 显示名缓存
119 }
120
121 type nameCacheEntry struct {
122 name string
123 expires time.Time
124 }
125
126 const (
127 userNameCacheTTL = time.Hour
128 userNameFallbackCacheTTL = 5 * time.Minute
129 )
130
131 // New 创建飞书 Bot 适配器。
132 func New(cfg config.FeishuBotConfig, logger *slog.Logger) bot.Adapter {
133 return &adapter{
134 cfg: cfg,
135 logger: logger.With("platform", "feishu"),
136 seen: make(map[string]bool),
137 }
138 }
139
140 func (a *adapter) Platform() bot.Platform { return bot.PlatformFeishu }
141 func (a *adapter) Name() string { return "feishu" }
142
143 func (a *adapter) Start(ctx context.Context) error {
144 a.msgCh = make(chan bot.InboundMessage, 64)
145 ctx, a.cancel = context.WithCancel(ctx)
146
147 mode := a.cfg.Mode
148 if mode == "" {
149 mode = "webhook"
150 }
151
152 switch mode {
153 case "webhook":
154 // Webhook mode exposes a public HTTP endpoint; without a verification
155 // token verificationTokenValid accepts every caller, so fail closed
156 // rather than let anyone drive the agent.
157 if strings.TrimSpace(a.cfg.VerificationToken) == "" {
158 return fmt.Errorf("feishu: webhook mode needs verification_token set — refusing to expose an unauthenticated event endpoint")
159 }
160 go a.runWebhook(ctx)
161 default:
162 if _, err := a.appSecret(); err != nil {
163 return err
164 }
165 go a.runWebSocket(ctx)
166 }
167 // bot open_id 用于把群聊 @ 门控收紧为“必须 @ 本 bot”;拉取失败只降级为
168 // 旧行为(任意 @ 放行),不阻塞启动。
169 go a.fetchBotOpenID(ctx)
170 return nil
171 }
172
173 func (a *adapter) botOpenID() string {
174 a.botMu.Lock()
175 defer a.botMu.Unlock()
176 return a.botID
177 }
178
179 func (a *adapter) fetchBotOpenID(ctx context.Context) {
180 client, err := a.sdkClient()
181 if err != nil {
182 return
183 }
184 ctx, cancel := context.WithTimeout(ctx, 15*time.Second)
185 defer cancel()
186 resp, err := client.Get(ctx, "/open-apis/bot/v3/info", nil, larkcore.AccessTokenTypeTenant)
187 if err != nil {
188 a.logger.Warn("feishu bot info fetch failed; group mention gating stays permissive", "err", err)
189 return
190 }
191 var payload struct {
192 Code int `json:"code"`
193 Bot struct {
194 OpenID string `json:"open_id"`
195 } `json:"bot"`
196 }
197 if err := json.Unmarshal(resp.RawBody, &payload); err != nil || payload.Code != 0 || payload.Bot.OpenID == "" {
198 a.logger.Warn("feishu bot info unavailable; group mention gating stays permissive", "code", payload.Code, "err", err)
199 return
200 }
201 a.botMu.Lock()
202 a.botID = payload.Bot.OpenID
203 a.botMu.Unlock()
204 a.logger.Info("feishu bot identity resolved", "open_id", logHash(payload.Bot.OpenID))
205 }
206
207 // resolveUserName 把 open_id 解析为显示名(1 小时缓存)。缺少 contact 权限或
208 // 调用失败时回退 open_id 本身,并短暂缓存回退值避免每条消息都打一次 API。
209 func (a *adapter) resolveUserName(ctx context.Context, openID string) string {
210 openID = strings.TrimSpace(openID)
211 if openID == "" {
212 return ""
213 }
214 now := time.Now()
215 a.nameMu.Lock()
216 if entry, ok := a.names[openID]; ok && now.Before(entry.expires) {
217 a.nameMu.Unlock()
218 return entry.name
219 }
220 a.nameMu.Unlock()
221 name, ttl := a.lookupUserName(ctx, openID)
222 a.nameMu.Lock()
223 if a.names == nil {
224 a.names = make(map[string]nameCacheEntry)
225 }
226 if len(a.names) > 10000 {
227 a.names = make(map[string]nameCacheEntry)
228 }
229 a.names[openID] = nameCacheEntry{name: name, expires: now.Add(ttl)}
230 a.nameMu.Unlock()
231 return name
232 }
233
234 func (a *adapter) lookupUserName(ctx context.Context, openID string) (string, time.Duration) {
235 client, err := a.sdkClient()
236 if err != nil {
237 return openID, userNameFallbackCacheTTL
238 }
239 ctx, cancel := context.WithTimeout(ctx, 5*time.Second)
240 defer cancel()
241 req := larkcontact.NewGetUserReqBuilder().
242 UserId(openID).
243 UserIdType(larkcontact.UserIdTypeOpenId).
244 Build()
245 resp, err := client.Contact.User.Get(ctx, req)
246 if err != nil || resp == nil || !resp.Success() || resp.Data == nil || resp.Data.User == nil {
247 return openID, userNameFallbackCacheTTL
248 }
249 name := stringPtrValue(resp.Data.User.Name)
250 if name == "" {
251 return openID, userNameFallbackCacheTTL
252 }
253 return name, userNameCacheTTL
254 }
255
256 func (a *adapter) Stop() error {
257 if a.cancel != nil {
258 a.cancel()
259 }
260 if a.wsClient != nil {
261 a.wsClient.Close()
262 }
263 return nil
264 }
265
266 func (a *adapter) Send(ctx context.Context, msg bot.OutboundMessage) (bot.SendResult, error) {
267 return a.sendMessage(ctx, msg)
268 }
269
270 func (a *adapter) SendTyping(ctx context.Context, chatID string) error {
271 return nil
272 }
273
274 func (a *adapter) Messages() <-chan bot.InboundMessage {
275 return a.msgCh
276 }
277
278 func (a *adapter) appSecret() (string, error) {
279 secret := os.Getenv(a.cfg.AppSecretEnv)
280 if a.cfg.AppID == "" || secret == "" {
281 return "", fmt.Errorf("feishu app_id or %s is not configured", a.cfg.AppSecretEnv)
282 }
283 return secret, nil
284 }
285
286 // runWebSocket 启动飞书 WebSocket 长连接。
287 func (a *adapter) runWebSocket(ctx context.Context) {
288 secret, err := a.appSecret()
289 if err != nil {
290 a.logger.Error("feishu websocket config error", "err", err)
291 return
292 }
293 eventHandler := a.newEventDispatcher()
294 bot.RunWithRetry(ctx, a.logger, "feishu sdk websocket", bot.RetryConfig{}, func(ctx context.Context) error {
295 opts := []larkws.ClientOption{
296 larkws.WithEventHandler(eventHandler),
297 larkws.WithLogLevel(larkcore.LogLevelError),
298 larkws.WithAutoReconnect(true),
299 larkws.WithOnReady(func() { a.logger.Info("feishu sdk websocket connected") }),
300 larkws.WithOnReconnecting(func() { a.logger.Warn("feishu sdk websocket reconnecting") }),
301 larkws.WithOnReconnected(func() { a.logger.Info("feishu sdk websocket reconnected") }),
302 larkws.WithOnError(func(err error) { a.logger.Error("feishu sdk websocket error", "err", err) }),
303 }
304 if feishuDomain(a.cfg.Domain) == "lark" {
305 opts = append(opts, larkws.WithDomain(lark.LarkBaseUrl))
306 }
307 client := larkws.NewClient(a.cfg.AppID, secret, opts...)
308 a.wsClient = client
309 // client.Start blocks; run it off-loop so cancellation closes the client
310 // immediately rather than waiting for Start to notice ctx. RunWithRetry
311 // handles the reconnect backoff.
312 errCh := make(chan error, 1)
313 go func() { errCh <- client.Start(ctx) }()
314 select {
315 case <-ctx.Done():
316 client.Close()
317 return nil
318 case err := <-errCh:
319 client.Close()
320 return err
321 }
322 })
323 }
324
325 func (a *adapter) newEventDispatcher() *dispatcher.EventDispatcher {
326 return dispatcher.NewEventDispatcher(a.cfg.VerificationToken, "").
327 OnP2MessageReceiveV1(func(ctx context.Context, event *larkim.P2MessageReceiveV1) error {
328 a.handleSDKMessage(ctx, event)
329 return nil
330 }).
331 OnP2MessageReadV1(func(ctx context.Context, event *larkim.P2MessageReadV1) error {
332 return nil
333 }).
334 OnP2MessageReactionCreatedV1(func(ctx context.Context, event *larkim.P2MessageReactionCreatedV1) error {
335 return nil
336 }).
337 OnP2MessageReactionDeletedV1(func(ctx context.Context, event *larkim.P2MessageReactionDeletedV1) error {
338 return nil
339 }).
340 OnP2CardActionTrigger(func(ctx context.Context, event *callback.CardActionTriggerEvent) (*callback.CardActionTriggerResponse, error) {
341 if event == nil || event.EventReq == nil || !a.handleCardAction(event.Body) {
342 a.logger.Warn("feishu card action ignored", "reason", "invalid_payload")
343 return cardActionToast("warning", "操作无效或已过期"), nil
344 }
345 return cardActionToast("success", "操作已提交"), nil
346 })
347 }
348
349 func (a *adapter) handleSDKMessage(ctx context.Context, event *larkim.P2MessageReceiveV1) {
350 if event == nil || event.Event == nil || event.Event.Message == nil {
351 return
352 }
353 eventID := ""
354 if event.EventV2Base != nil && event.EventV2Base.Header != nil {
355 eventID = event.EventV2Base.Header.EventID
356 }
357 if eventID != "" {
358 if a.markSeen(eventID) {
359 return
360 }
361 }
362 msg := event.Event.Message
363 messageID := stringPtrValue(msg.MessageId)
364 mentions := sdkMentionRefs(msg.Mentions)
365 chatType := bot.ChatDM
366 if stringPtrValue(msg.ChatType) == "group" || stringPtrValue(msg.ChatType) == "topic_group" {
367 chatType = bot.ChatGroup
368 if a.cfg.RequireMention && !a.mentionsBot(mentions) {
369 a.logger.Info("feishu message ignored", "reason", "missing_mention", "chat", logHash(stringPtrValue(msg.ChatId)), "message", logHash(messageID))
370 return
371 }
372 }
373 msgType := stringPtrValue(msg.MessageType)
374 text, media, ok := a.parseInboundContent(msgType, stringPtrValue(msg.Content), messageID)
375 if !ok {
376 a.logger.Info("feishu message ignored", "reason", "unsupported_type", "msg_type", msgType, "chat_type", stringPtrValue(msg.ChatType), "message", logHash(messageID))
377 return
378 }
379 text = a.replaceMentionPlaceholders(text, mentions)
380 if strings.TrimSpace(text) == "" && len(media) == 0 {
381 a.logger.Info("feishu message ignored", "reason", "empty_after_parse", "msg_type", msgType, "message", logHash(messageID))
382 return
383 }
384 userID := ""
385 senderOpenID := ""
386 if event.Event.Sender != nil && event.Event.Sender.SenderId != nil {
387 senderOpenID = stringPtrValue(event.Event.Sender.SenderId.OpenId)
388 userID = firstNonEmpty(
389 senderOpenID,
390 stringPtrValue(event.Event.Sender.SenderId.UnionId),
391 stringPtrValue(event.Event.Sender.SenderId.UserId),
392 )
393 }
394 userName := userID
395 var resolveUserName func(context.Context) string
396 if senderOpenID != "" {
397 resolveUserName = func(ctx context.Context) string {
398 return a.resolveUserName(ctx, senderOpenID)
399 }
400 }
401 ib := bot.InboundMessage{
402 Platform: bot.PlatformFeishu,
403 ChatType: chatType,
404 ChatID: stringPtrValue(msg.ChatId),
405 UserID: userID,
406 UserName: userName,
407 Text: text,
408 MessageID: messageID,
409 ThreadID: stringPtrValue(msg.ThreadId),
410 Media: media,
411 ResolveUserName: resolveUserName,
412 Raw: event,
413 }
414 select {
415 case a.msgCh <- ib:
416 a.logger.Info("feishu inbound queued", "chat_type", chatType, "msg_type", msgType, "chat", logHash(ib.ChatID), "user", logHash(ib.UserID), "message", logHash(ib.MessageID), "text_chars", len([]rune(ib.Text)), "media_items", len(media))
417 default:
418 a.logger.Warn("feishu message channel full")
419 }
420 }
421
422 func (a *adapter) handleWSEvent(ctx context.Context, raw json.RawMessage) {
423 var evt feishuEvent
424 if err := json.Unmarshal(raw, &evt); err != nil {
425 return
426 }
427
428 if a.markSeen(evt.Header.EventID) {
429 return
430 }
431
432 switch evt.Header.EventType {
433 case "im.message.receive_v1":
434 var msg feishuMsgEvent
435 if err := json.Unmarshal(evt.Event, &msg); err != nil {
436 return
437 }
438 a.handleMessage(ctx, msg)
439 }
440 }
441
442 func (a *adapter) handleCardAction(raw []byte) bool {
443 var payload struct {
444 Header feishuHeader `json:"header"`
445 Event struct {
446 Operator struct {
447 UserID string `json:"user_id"`
448 OpenID string `json:"open_id"`
449 UnionID string `json:"union_id"`
450 OperatorID struct {
451 UserID string `json:"user_id"`
452 OpenID string `json:"open_id"`
453 UnionID string `json:"union_id"`
454 } `json:"operator_id"`
455 } `json:"operator"`
456 Context struct {
457 OpenMessageID string `json:"open_message_id"`
458 OpenChatID string `json:"open_chat_id"`
459 } `json:"context"`
460 Action struct {
461 Value map[string]string `json:"value"`
462 } `json:"action"`
463 } `json:"event"`
464 }
465 if err := json.Unmarshal(raw, &payload); err != nil {
466 return false
467 }
468 command := payload.Event.Action.Value["command"]
469 if command == "" || payload.Event.Context.OpenChatID == "" {
470 return false
471 }
472 if a.markSeen(payload.Header.EventID) {
473 return true
474 }
475 chatType := cardActionChatType(payload.Event.Action.Value["chat_type"])
476 operatorID := firstNonEmpty(
477 payload.Event.Operator.OperatorID.UnionID,
478 payload.Event.Operator.OperatorID.OpenID,
479 payload.Event.Operator.OperatorID.UserID,
480 payload.Event.Operator.UnionID,
481 payload.Event.Operator.OpenID,
482 payload.Event.Operator.UserID,
483 )
484 routeUserID := firstNonEmpty(payload.Event.Action.Value["user_id"], operatorID)
485 ib := bot.InboundMessage{
486 Platform: bot.PlatformFeishu,
487 ChatType: chatType,
488 ChatID: payload.Event.Context.OpenChatID,
489 UserID: routeUserID,
490 UserName: routeUserID,
491 OperatorID: operatorID,
492 Text: command,
493 MessageID: payload.Event.Context.OpenMessageID,
494 }
495 select {
496 case a.msgCh <- ib:
497 default:
498 a.logger.Warn("feishu card action channel full")
499 }
500 return true
501 }
502
503 func (a *adapter) markSeen(eventID string) bool {
504 if eventID == "" {
505 return false
506 }
507 a.seenMu.Lock()
508 defer a.seenMu.Unlock()
509 if a.seen == nil {
510 a.seen = make(map[string]bool)
511 }
512 if a.seen[eventID] {
513 return true
514 }
515 a.seen[eventID] = true
516 if len(a.seen) > 10000 {
517 a.seen = make(map[string]bool)
518 a.seen[eventID] = true
519 }
520 return false
521 }
522
523 func cardActionChatType(raw string) bot.ChatType {
524 switch bot.ChatType(raw) {
525 case bot.ChatDM, bot.ChatGroup, bot.ChatGuild, bot.ChatDirect, bot.ChatThread:
526 return bot.ChatType(raw)
527 default:
528 return bot.ChatGroup
529 }
530 }
531
532 func cardActionToast(toastType, content string) *callback.CardActionTriggerResponse {
533 return &callback.CardActionTriggerResponse{
534 Toast: &callback.Toast{
535 Type: toastType,
536 Content: content,
537 },
538 }
539 }
540
541 func (a *adapter) verificationTokenValid(token string) bool {
542 if a.cfg.VerificationToken == "" {
543 return false
544 }
545 return subtle.ConstantTimeCompare([]byte(token), []byte(a.cfg.VerificationToken)) == 1
546 }
547
548 func firstNonEmpty(vals ...string) string {
549 for _, v := range vals {
550 if v != "" {
551 return v
552 }
553 }
554 return ""
555 }
556
557 func logHash(id string) string {
558 if id == "" {
559 return ""
560 }
561 sum := sha256.Sum256([]byte(id))
562 return hex.EncodeToString(sum[:])[:12]
563 }
564
565 func (a *adapter) handleMessage(ctx context.Context, msg feishuMsgEvent) {
566 mentions := webhookMentionRefs(msg.Mentions)
567
568 // @mention gating:仅在群聊中检查是否 @了 bot
569 chatType := bot.ChatDM
570 if msg.ChatType == "group" || msg.ChatType == "topic_group" {
571 chatType = bot.ChatGroup
572 if a.cfg.RequireMention && !a.mentionsBot(mentions) {
573 a.logger.Info("feishu message ignored", "reason", "missing_mention", "chat", logHash(msg.ChatID), "message", logHash(msg.MessageID))
574 return
575 }
576 }
577
578 text, media, ok := a.parseInboundContent(msg.MsgType, msg.Content, msg.MessageID)
579 if !ok {
580 a.logger.Info("feishu message ignored", "reason", "unsupported_type", "msg_type", msg.MsgType, "chat_type", msg.ChatType, "message", logHash(msg.MessageID))
581 return
582 }
583 text = a.replaceMentionPlaceholders(text, mentions)
584 if strings.TrimSpace(text) == "" && len(media) == 0 {
585 a.logger.Info("feishu message ignored", "reason", "empty_after_parse", "msg_type", msg.MsgType, "message", logHash(msg.MessageID))
586 return
587 }
588
589 userName := msg.Sender.SenderID.OpenID
590 var resolveUserName func(context.Context) string
591 if userName != "" {
592 openID := msg.Sender.SenderID.OpenID
593 resolveUserName = func(ctx context.Context) string {
594 return a.resolveUserName(ctx, openID)
595 }
596 }
597 ib := bot.InboundMessage{
598 Platform: bot.PlatformFeishu,
599 ChatType: chatType,
600 ChatID: msg.ChatID,
601 UserID: msg.Sender.SenderID.OpenID,
602 UserName: userName,
603 Text: text,
604 MessageID: msg.MessageID,
605 ThreadID: msg.ThreadID,
606 Media: media,
607 ResolveUserName: resolveUserName,
608 }
609
610 select {
611 case a.msgCh <- ib:
612 a.logger.Info("feishu inbound queued", "chat_type", chatType, "msg_type", msg.MsgType, "chat", logHash(ib.ChatID), "user", logHash(ib.UserID), "message", logHash(ib.MessageID), "text_chars", len([]rune(ib.Text)), "media_items", len(media))
613 default:
614 a.logger.Warn("feishu message channel full")
615 }
616 }
617
618 // SendText sends an interactive card with markdown content to a Feishu/Lark chat_id using the SDK.
619 // It is used by the desktop settings panel as an actual connection test.
620 func SendText(ctx context.Context, cfg config.FeishuBotConfig, chatID, text string) (bot.SendResult, error) {
621 a := &adapter{cfg: cfg, logger: slog.Default().With("platform", "feishu")}
622 return a.sendMessage(ctx, bot.OutboundMessage{ChatID: chatID, Text: text})
623 }
624
625 // sendMessage 使用飞书/Lark SDK 以 Interactive Card (JSON 2.0) 发送消息。
626 // Card 内嵌 markdown 元素,支持 CommonMark 标准语法。
627 // 当卡片体积超过 30KB 限制(如大段代码),自动降级为纯文本消息。
628 // MediaURLs are bare filenames staged in an operator-configured outbound media
629 // root. URL fetching and arbitrary-path reads are intentionally unsupported.
630 func (a *adapter) sendMessage(ctx context.Context, msg bot.OutboundMessage) (bot.SendResult, error) {
631 if msg.Card != nil {
632 return a.sendCard(ctx, msg)
633 }
634 if len(msg.MediaURLs) == 0 {
635 return a.sendRenderedText(ctx, msg)
636 }
637 media, err := a.loadOutboundMedia(msg.MediaURLs)
638 if err != nil {
639 return bot.SendResult{}, err
640 }
641
642 var result bot.SendResult
643 if strings.TrimSpace(msg.Text) != "" {
644 textResult, err := a.sendRenderedText(ctx, msg)
645 result.Merge(textResult)
646 if err != nil {
647 return result, err
648 }
649 }
650 mediaResult, err := a.sendMedia(ctx, msg, media)
651 result.Merge(mediaResult)
652 return result, err
653 }
654
655 func (a *adapter) sendRenderedText(ctx context.Context, msg bot.OutboundMessage) (bot.SendResult, error) {
656 cardContent, err := buildMarkdownCard(msg.Text)
657 if err != nil {
658 a.logger.Warn("build markdown card failed, falling back to text", "err", err)
659 return a.sendSDKContent(ctx, msg, larkim.MsgTypeText, feishuTextContent(msg.Text))
660 }
661 result, err := a.sendSDKContent(ctx, msg, larkim.MsgTypeInteractive, cardContent)
662 if err != nil && isCardLimitError(err) {
663 a.logger.Warn("card send failed (size limit), retrying as text", "err", err)
664 return a.sendSDKContent(ctx, msg, larkim.MsgTypeText, feishuTextContent(msg.Text))
665 }
666 return result, err
667 }
668
669 func buildMarkdownCard(content string) (string, error) {
670 card := map[string]any{
671 "schema": "2.0",
672 // update_multi marks the card as a shared card that can be patched for
673 // all recipients after sending; without it Im.Message.Patch (used by
674 // EditMessage for streaming) is rejected, which would collapse
675 // streaming into a flood of new messages. See references/desktop-ui.
676 "config": map[string]any{
677 "update_multi": true,
678 },
679 "body": map[string]any{
680 "elements": []map[string]any{
681 {
682 "tag": "markdown",
683 "content": content,
684 },
685 },
686 },
687 }
688 data, err := json.Marshal(card)
689 if err != nil {
690 return "", err
691 }
692 return string(data), nil
693 }
694
695 func feishuTextContent(text string) string {
696 content, _ := json.Marshal(textContent{Text: text})
697 return string(content)
698 }
699
700 func isCardLimitError(err error) bool {
701 if err == nil {
702 return false
703 }
704 s := err.Error()
705 return strings.Contains(s, "11310") || strings.Contains(s, "11325")
706 }
707
708 const feishuReplyRecalledCode = 230011
709
710 type feishuAPIError struct {
711 op string
712 code int
713 msg string
714 }
715
716 func (e *feishuAPIError) Error() string {
717 return fmt.Sprintf("feishu %s error: %s", e.op, feishuCodeError(e.code, e.msg))
718 }
719
720 func isReplyFallbackError(err error) bool {
721 var apiErr *feishuAPIError
722 return errors.As(err, &apiErr) && apiErr.op == "reply" && apiErr.code == feishuReplyRecalledCode
723 }
724
725 // sdkClient lazily builds the shared lark client. It is called concurrently —
726 // the fetchBotOpenID goroutine, per-message resolveUserName, and per-resource
727 // downloads all race on first use at startup — so the check-and-build is guarded
728 // by clientMu (a bare a.client read/write would data-race, tripping -race).
729 func (a *adapter) sdkClient() (*lark.Client, error) {
730 a.clientMu.Lock()
731 defer a.clientMu.Unlock()
732 if a.client != nil {
733 return a.client, nil
734 }
735 secret, err := a.appSecret()
736 if err != nil {
737 return nil, err
738 }
739 opts := []lark.ClientOptionFunc{
740 lark.WithLogLevel(larkcore.LogLevelError),
741 lark.WithReqTimeout(15 * time.Second),
742 lark.WithSource("reasonix"),
743 }
744 if feishuDomain(a.cfg.Domain) == "lark" {
745 opts = append(opts, lark.WithOpenBaseUrl(lark.LarkBaseUrl), lark.WithOAuthBaseUrl(lark.OAuthBaseUrlLark))
746 }
747 a.client = lark.NewClient(a.cfg.AppID, secret, opts...)
748 return a.client, nil
749 }
750
751 func (a *adapter) sendSDKContent(ctx context.Context, msg bot.OutboundMessage, msgType, content string) (bot.SendResult, error) {
752 client, err := a.sdkClient()
753 if err != nil {
754 return bot.SendResult{}, err
755 }
756 chatID := strings.TrimSpace(msg.ChatID)
757 if chatID == "" {
758 return bot.SendResult{}, fmt.Errorf("feishu chat_id is empty")
759 }
760 // 带触发消息 ID 时用 Reply 引用回复:话题群里回复会落到对应话题,
761 // 普通群里带引用上下文。只有飞书明确返回“消息已撤回”时才回退普通
762 // 发送;传输错误的提交结果不确定,回退 Create 可能产生重复消息。
763 if replyTo := strings.TrimSpace(msg.ReplyToMsgID); replyTo != "" {
764 result, err := a.replySDKContent(ctx, replyTo, msgType, content)
765 if err == nil {
766 return result, nil
767 }
768 if !isReplyFallbackError(err) {
769 return bot.SendResult{}, err
770 }
771 a.logger.Warn("feishu reply failed; falling back to create", "message", logHash(replyTo), "err", err)
772 }
773 // Stable across retries so a retry after a post-commit connection drop does
774 // not send a duplicate visible message (Feishu dedups on uuid).
775 uuid := newIdempotencyKey()
776 var result bot.SendResult
777 err = withTransientRetry(ctx, a.logger, "create message", func(ctx context.Context) error {
778 body := larkim.NewCreateMessageReqBodyBuilder().ReceiveId(chatID).MsgType(msgType).Content(content)
779 if uuid != "" {
780 body = body.Uuid(uuid)
781 }
782 req := larkim.NewCreateMessageReqBuilder().
783 ReceiveIdType(larkim.CreateMessageV1ReceiveIDTypeChatId).
784 Body(body.Build()).
785 Build()
786 resp, err := client.Im.Message.Create(ctx, req)
787 if err != nil {
788 return err
789 }
790 if resp == nil {
791 return fmt.Errorf("feishu send error: empty response")
792 }
793 if !resp.Success() {
794 return fmt.Errorf("feishu send error: %s", feishuCodeError(resp.Code, resp.Msg))
795 }
796 if resp.Data != nil {
797 result = bot.SendResult{MessageID: stringPtrValue(resp.Data.MessageId)}
798 }
799 return nil
800 })
801 if err != nil {
802 return bot.SendResult{}, err
803 }
804 return result, nil
805 }
806
807 func (a *adapter) replySDKContent(ctx context.Context, replyTo, msgType, content string) (bot.SendResult, error) {
808 client, err := a.sdkClient()
809 if err != nil {
810 return bot.SendResult{}, err
811 }
812 uuid := newIdempotencyKey()
813 var result bot.SendResult
814 err = withTransientRetry(ctx, a.logger, "reply message", func(ctx context.Context) error {
815 body := larkim.NewReplyMessageReqBodyBuilder().MsgType(msgType).Content(content)
816 if uuid != "" {
817 body = body.Uuid(uuid)
818 }
819 req := larkim.NewReplyMessageReqBuilder().
820 MessageId(replyTo).
821 Body(body.Build()).
822 Build()
823 resp, err := client.Im.Message.Reply(ctx, req)
824 if err != nil {
825 return err
826 }
827 if resp == nil {
828 return fmt.Errorf("feishu reply error: empty response")
829 }
830 if !resp.Success() {
831 return &feishuAPIError{op: "reply", code: resp.Code, msg: resp.Msg}
832 }
833 if resp.Data != nil {
834 result = bot.SendResult{MessageID: stringPtrValue(resp.Data.MessageId)}
835 }
836 return nil
837 })
838 if err != nil {
839 return bot.SendResult{}, err
840 }
841 return result, nil
842 }
843
844 func (a *adapter) AddPendingReaction(ctx context.Context, messageID string) (func(), error) {
845 messageID = strings.TrimSpace(messageID)
846 if messageID == "" {
847 return nil, nil
848 }
849 client, err := a.sdkClient()
850 if err != nil {
851 return nil, err
852 }
853 req := larkim.NewCreateMessageReactionReqBuilder().
854 MessageId(messageID).
855 Body(larkim.NewCreateMessageReactionReqBodyBuilder().
856 ReactionType(larkim.NewEmojiBuilder().EmojiType(feishuPendingReactionEmoji).Build()).
857 Build()).
858 Build()
859 resp, err := client.Im.MessageReaction.Create(ctx, req)
860 if err != nil {
861 return nil, err
862 }
863 if resp == nil || !resp.Success() {
864 if resp != nil {
865 return nil, fmt.Errorf("feishu reaction error: %s", feishuCodeError(resp.Code, resp.Msg))
866 }
867 return nil, fmt.Errorf("feishu reaction error: empty response")
868 }
869 reactionID := ""
870 if resp.Data != nil && resp.Data.ReactionId != nil {
871 reactionID = *resp.Data.ReactionId
872 }
873 if reactionID == "" {
874 return nil, nil
875 }
876 cleanup := func() {
877 delReq := larkim.NewDeleteMessageReactionReqBuilder().
878 MessageId(messageID).
879 ReactionId(reactionID).
880 Build()
881 if _, err := client.Im.MessageReaction.Delete(context.Background(), delReq); err != nil {
882 a.logger.Warn("feishu reaction cleanup failed", "message", logHash(messageID), "err", err)
883 }
884 }
885 return cleanup, nil
886 }
887
888 // sendCard 发送 interactive card 消息(用于审批/问答)。
889 func (a *adapter) sendCard(ctx context.Context, msg bot.OutboundMessage) (bot.SendResult, error) {
890 card := msg.Card
891
892 elements := make([]map[string]any, 0)
893 for _, el := range card.Elements {
894 item := map[string]any{"tag": el.Tag}
895 if el.Content != "" {
896 item["content"] = el.Content
897 }
898 if actions, ok := el.Extra["actions"]; ok && el.Tag == "action" {
899 item["actions"] = actions
900 } else {
901 maps.Copy(item, el.Extra)
902 }
903 elements = append(elements, item)
904 }
905
906 cardPayload := map[string]any{
907 "header": map[string]any{
908 "title": map[string]string{
909 "tag": "plain_text",
910 "content": card.Header,
911 },
912 },
913 "elements": elements,
914 }
915
916 cardJSON, _ := json.Marshal(cardPayload)
917 return a.sendSDKContent(ctx, msg, larkim.MsgTypeInteractive, string(cardJSON))
918 }
919
920 func feishuDomain(domain string) string {
921 if strings.EqualFold(strings.TrimSpace(domain), "lark") {
922 return "lark"
923 }
924 return "feishu"
925 }
926
927 func stringPtrValue(ptr *string) string {
928 if ptr == nil {
929 return ""
930 }
931 return strings.TrimSpace(*ptr)
932 }
933
934 func feishuCodeError(code int, msg string) string {
935 msg = strings.TrimSpace(msg)
936 if msg == "" {
937 msg = "unknown error"
938 }
939 if code == 0 {
940 return msg
941 }
942 return fmt.Sprintf("%s (code %d)", msg, code)
943 }
944
945 // runWebhook 启动飞书 Webhook 模式。
946 func (a *adapter) runWebhook(ctx context.Context) {
947 port := a.cfg.WebhookPort
948 if port == 0 {
949 port = 8080
950 }
951
952 mux := http.NewServeMux()
953 mux.HandleFunc("/feishu/event", func(w http.ResponseWriter, r *http.Request) {
954 body, err := io.ReadAll(io.LimitReader(r.Body, 1024*1024))
955 if err != nil {
956 http.Error(w, "bad request", http.StatusBadRequest)
957 return
958 }
959 var challenge struct {
960 Challenge string `json:"challenge"`
961 Token string `json:"token"`
962 Type string `json:"type"`
963 }
964 _ = json.Unmarshal(body, &challenge)
965 if challenge.Type == "url_verification" {
966 if !a.verificationTokenValid(challenge.Token) {
967 http.Error(w, "forbidden", http.StatusForbidden)
968 return
969 }
970 w.Header().Set("Content-Type", "application/json")
971 if err := json.NewEncoder(w).Encode(map[string]string{"challenge": challenge.Challenge}); err != nil {
972 a.logger.Error("feishu challenge response error", "err", err)
973 }
974 return
975 }
976
977 var evt feishuEvent
978 if err := json.Unmarshal(body, &evt); err != nil {
979 http.Error(w, "bad request", http.StatusBadRequest)
980 return
981 }
982 if !a.verificationTokenValid(evt.Header.Token) {
983 http.Error(w, "forbidden", http.StatusForbidden)
984 return
985 }
986
987 if !a.handleCardAction(body) {
988 raw, _ := json.Marshal(evt)
989 a.handleWSEvent(ctx, raw)
990 }
991 w.WriteHeader(http.StatusOK)
992 })
993
994 server := &http.Server{
995 Addr: fmt.Sprintf(":%d", port),
996 Handler: mux,
997 }
998
999 go func() {
1000 <-ctx.Done()
1001 if err := server.Shutdown(context.Background()); err != nil && !errors.Is(err, http.ErrServerClosed) {
1002 a.logger.Error("feishu webhook shutdown error", "err", err)
1003 }
1004 }()
1005
1006 a.logger.Info("feishu webhook listening", "port", port)
1007 if err := server.ListenAndServe(); err != http.ErrServerClosed {
1008 a.logger.Error("feishu webhook server error", "err", err)
1009 }
1010 }
1011
1011 lines GO