Delete internal directory

This commit is contained in:
公明
2026-08-15 01:54:26 +08:00
committed by GitHub
parent 5ce66ee2f8
commit 2910289f0b
591 changed files with 0 additions and 136567 deletions
-6
View File
@@ -1,6 +0,0 @@
package robot
// MessageHandler 供飞书/钉钉长连接调用的消息处理接口(由 handler.RobotHandler 实现)
type MessageHandler interface {
HandleMessage(platform, userID, text string) string
}
-151
View File
@@ -1,151 +0,0 @@
package robot
import (
"bytes"
"context"
"encoding/json"
"net/http"
"strings"
"time"
"cyberstrike-ai/internal/config"
"github.com/open-dingtalk/dingtalk-stream-sdk-go/chatbot"
"github.com/open-dingtalk/dingtalk-stream-sdk-go/client"
dingutils "github.com/open-dingtalk/dingtalk-stream-sdk-go/utils"
"go.uber.org/zap"
)
const (
dingReconnectInitial = 5 * time.Second // 首次重连间隔
dingReconnectMax = 60 * time.Second // 最大重连间隔
)
// StartDing 启动钉钉 Stream 长连接(无需公网),收到消息后调用 handler 并通过 SessionWebhook 回复。
// 断线(如笔记本睡眠、网络中断)后会自动重连;ctx 被取消时退出,便于配置变更时重启。
func StartDing(ctx context.Context, robotsCfg config.RobotsConfig, h MessageHandler, logger *zap.Logger) {
cfg := robotsCfg.Dingtalk
if !cfg.Enabled || cfg.ClientID == "" || cfg.ClientSecret == "" {
return
}
go runDingLoop(ctx, cfg, robotsCfg.Session.StrictUserIdentityEnabled(), h, logger)
}
// runDingLoop 循环维持钉钉长连接:断开且 ctx 未取消时按退避间隔重连。
func runDingLoop(ctx context.Context, cfg config.RobotDingtalkConfig, strictUserIdentity bool, h MessageHandler, logger *zap.Logger) {
backoff := dingReconnectInitial
for {
streamClient := client.NewStreamClient(
client.WithAppCredential(client.NewAppCredentialConfig(cfg.ClientID, cfg.ClientSecret)),
client.WithSubscription(dingutils.SubscriptionTypeKCallback, "/v1.0/im/bot/messages/get",
chatbot.NewDefaultChatBotFrameHandler(func(ctx context.Context, msg *chatbot.BotCallbackDataModel) ([]byte, error) {
go handleDingMessage(ctx, msg, cfg, strictUserIdentity, h, logger)
return nil, nil
}).OnEventReceived),
)
logger.Info("钉钉 Stream 正在连接…", zap.String("client_id", cfg.ClientID))
err := streamClient.Start(ctx)
if ctx.Err() != nil {
logger.Info("钉钉 Stream 已按配置重启关闭")
return
}
if err != nil {
logger.Warn("钉钉 Stream 长连接断开(如睡眠/断网),将自动重连", zap.Error(err), zap.Duration("retry_after", backoff))
}
select {
case <-ctx.Done():
return
case <-time.After(backoff):
// 下次重连间隔递增,上限 60 秒,避免频繁重试
if backoff < dingReconnectMax {
backoff *= 2
if backoff > dingReconnectMax {
backoff = dingReconnectMax
}
}
}
}
}
func handleDingMessage(ctx context.Context, msg *chatbot.BotCallbackDataModel, cfg config.RobotDingtalkConfig, strictUserIdentity bool, h MessageHandler, logger *zap.Logger) {
if msg == nil || msg.SessionWebhook == "" {
return
}
content := ""
if msg.Text.Content != "" {
content = strings.TrimSpace(msg.Text.Content)
}
if content == "" && msg.Msgtype == "richText" {
if cMap, ok := msg.Content.(map[string]interface{}); ok {
if rich, ok := cMap["richText"].([]interface{}); ok {
for _, c := range rich {
if m, ok := c.(map[string]interface{}); ok {
if txt, ok := m["text"].(string); ok {
content = strings.TrimSpace(txt)
break
}
}
}
}
}
}
if content == "" {
logger.Debug("钉钉消息内容为空,已忽略", zap.String("msgtype", msg.Msgtype))
return
}
logger.Info("钉钉收到消息", zap.String("sender", msg.SenderId), zap.String("content", content))
tenantKey := strings.TrimSpace(cfg.ClientID)
if tenantKey == "" {
tenantKey = "default"
}
userID := strings.TrimSpace(msg.SenderId)
if userID != "" {
userID = "t:" + tenantKey + "|u:" + userID
} else if cfg.AllowConversationIDFallback && !strictUserIdentity {
conversationID := strings.TrimSpace(msg.ConversationId)
if conversationID != "" {
userID = "t:" + tenantKey + "|c:" + conversationID
}
}
if userID == "" {
logger.Warn("钉钉消息缺少可用用户标识,已忽略")
return
}
reply := h.HandleMessage("dingtalk", userID, content)
// 使用 markdown 类型以便正确展示标题、列表、代码块等格式
title := reply
if idx := strings.IndexAny(reply, "\n"); idx > 0 {
title = strings.TrimSpace(reply[:idx])
}
if len(title) > 50 {
title = title[:50] + "…"
}
if title == "" {
title = "回复"
}
body := map[string]interface{}{
"msgtype": "markdown",
"markdown": map[string]string{
"title": title,
"text": reply,
},
}
bodyBytes, _ := json.Marshal(body)
req, err := http.NewRequestWithContext(ctx, http.MethodPost, msg.SessionWebhook, bytes.NewReader(bodyBytes))
if err != nil {
logger.Warn("钉钉构造回复请求失败", zap.Error(err))
return
}
req.Header.Set("Content-Type", "application/json")
resp, err := http.DefaultClient.Do(req)
if err != nil {
logger.Warn("钉钉回复请求失败", zap.Error(err))
return
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
logger.Warn("钉钉回复非 200", zap.Int("status", resp.StatusCode))
return
}
logger.Debug("钉钉回复成功", zap.String("content_preview", reply))
}
-121
View File
@@ -1,121 +0,0 @@
package robot
import (
"context"
"strings"
"cyberstrike-ai/internal/config"
"github.com/bwmarrin/discordgo"
"go.uber.org/zap"
)
const (
discordPlatform = "discord"
discordMaxMessageRunes = 2000
)
// StartDiscord 启动 Discord GatewayWebSocket,无需公网回调)。
func StartDiscord(ctx context.Context, robotsCfg config.RobotsConfig, h MessageHandler, logger *zap.Logger) {
cfg := robotsCfg.Discord
if !cfg.Enabled || strings.TrimSpace(cfg.BotToken) == "" {
return
}
go runDiscordLoop(ctx, cfg, h, logger)
}
func runDiscordLoop(ctx context.Context, cfg config.RobotDiscordConfig, h MessageHandler, logger *zap.Logger) {
backoff := reconnectInitial
for {
err := runDiscordSession(ctx, cfg, h, logger)
if ctx.Err() != nil {
logger.Info("Discord Gateway 已按配置关闭")
return
}
if err != nil {
logger.Warn("Discord Gateway 异常,将自动重连", zap.Error(err), zap.Duration("retry_after", backoff))
}
if !waitReconnect(ctx, &backoff) {
return
}
}
}
func runDiscordSession(ctx context.Context, cfg config.RobotDiscordConfig, h MessageHandler, logger *zap.Logger) error {
token := strings.TrimSpace(cfg.BotToken)
if !strings.HasPrefix(token, "Bot ") {
token = "Bot " + token
}
session, err := discordgo.New(token)
if err != nil {
return err
}
session.Identify.Intents = discordgo.IntentsGuildMessages |
discordgo.IntentsDirectMessages |
discordgo.IntentMessageContent
session.AddHandler(func(s *discordgo.Session, m *discordgo.MessageCreate) {
if m == nil || m.Author == nil || m.Author.Bot {
return
}
text := strings.TrimSpace(m.Content)
if text == "" {
return
}
if m.GuildID != "" {
if !cfg.AllowGuildMessages {
return
}
if s.State.User == nil || !discordMentionsBot(m, s.State.User.ID) {
return
}
}
userID := discordSessionKey(m.GuildID, m.Author.ID)
logger.Info("Discord 收到消息", zap.String("from", userID), zap.String("content", text))
reply := h.HandleMessage(discordPlatform, userID, text)
discordPostReply(s, m.ChannelID, reply, logger)
})
if err := session.Open(); err != nil {
return err
}
logger.Info("Discord Gateway 已连接,等待收消息")
defer session.Close()
<-ctx.Done()
return ctx.Err()
}
func discordMentionsBot(m *discordgo.MessageCreate, botUserID string) bool {
if m == nil || botUserID == "" {
return false
}
for _, mention := range m.Mentions {
if mention != nil && mention.ID == botUserID {
return true
}
}
return strings.Contains(m.Content, "<@"+botUserID+">") || strings.Contains(m.Content, "<@!"+botUserID+">")
}
func discordSessionKey(guildID, userID string) string {
guildID = strings.TrimSpace(guildID)
userID = strings.TrimSpace(userID)
if guildID == "" {
return "u:" + userID
}
return "g:" + guildID + "|u:" + userID
}
func discordPostReply(s *discordgo.Session, channelID, reply string, logger *zap.Logger) {
reply = trimReply(reply)
if reply == "" {
return
}
for _, chunk := range splitTextChunks(reply, discordMaxMessageRunes) {
if _, err := s.ChannelMessageSend(channelID, chunk); err != nil {
logger.Warn("Discord 发送回复失败", zap.String("channel", channelID), zap.Error(err))
return
}
}
}
-316
View File
@@ -1,316 +0,0 @@
package ilink
import (
"bytes"
"context"
"crypto/rand"
"encoding/base64"
"encoding/json"
"fmt"
"io"
"net/http"
"net/url"
"strconv"
"strings"
"time"
)
const (
DefaultBaseURL = "https://ilinkai.weixin.qq.com"
DefaultBotType = "3"
DefaultBotAgent = "CyberStrikeAI/1.0"
ILinkAppID = "bot"
QRLongPollTimeout = 35 * time.Second
APIDefaultTimeout = 15 * time.Second
GetUpdatesTimeout = 35 * time.Second
)
// Client 微信 iLink Bot HTTP 客户端(与 @tencent-weixin/openclaw-weixin 协议兼容)
type Client struct {
BaseURL string
BotToken string
BotAgent string
ClientVersion uint32
HTTP *http.Client
}
func NewClient(baseURL, botToken, botAgent string, clientVersion uint32) *Client {
base := strings.TrimSpace(baseURL)
if base == "" {
base = DefaultBaseURL
}
agent := strings.TrimSpace(botAgent)
if agent == "" {
agent = DefaultBotAgent
}
return &Client{
BaseURL: strings.TrimRight(base, "/"),
BotToken: strings.TrimSpace(botToken),
BotAgent: sanitizeBotAgent(agent),
ClientVersion: clientVersion,
HTTP: &http.Client{Timeout: 0},
}
}
// BuildClientVersion 将 semver 编码为 iLink-App-ClientVersion0x00MMNNPP
func BuildClientVersion(version string) uint32 {
parts := strings.Split(version, ".")
parse := func(i int) int {
if i >= len(parts) {
return 0
}
n, _ := strconv.Atoi(strings.TrimSpace(parts[i]))
if n < 0 {
return 0
}
return n
}
major := parse(0) & 0xff
minor := parse(1) & 0xff
patch := parse(2) & 0xff
return uint32((major << 16) | (minor << 8) | patch)
}
type baseInfo struct {
ChannelVersion string `json:"channel_version"`
BotAgent string `json:"bot_agent"`
}
func (c *Client) buildBaseInfo() baseInfo {
return baseInfo{
ChannelVersion: "1.0.0",
BotAgent: c.BotAgent,
}
}
func randomWechatUIN() string {
var b [4]byte
_, _ = rand.Read(b[:])
u := uint32(b[0])<<24 | uint32(b[1])<<16 | uint32(b[2])<<8 | uint32(b[3])
return base64.StdEncoding.EncodeToString([]byte(strconv.FormatUint(uint64(u), 10)))
}
func (c *Client) commonHeaders() http.Header {
h := http.Header{}
h.Set("iLink-App-Id", ILinkAppID)
h.Set("iLink-App-ClientVersion", strconv.FormatUint(uint64(c.ClientVersion), 10))
return h
}
func (c *Client) authHeaders() http.Header {
h := c.commonHeaders()
h.Set("Content-Type", "application/json")
h.Set("AuthorizationType", "ilink_bot_token")
h.Set("X-WECHAT-UIN", randomWechatUIN())
if c.BotToken != "" {
h.Set("Authorization", "Bearer "+c.BotToken)
}
return h
}
func (c *Client) endpointURL(path string) (string, error) {
u, err := url.Parse(c.BaseURL + "/")
if err != nil {
return "", err
}
ref, err := url.Parse(path)
if err != nil {
return "", err
}
return u.ResolveReference(ref).String(), nil
}
func (c *Client) doRequest(ctx context.Context, method, path string, body []byte, headers http.Header, timeout time.Duration) ([]byte, error) {
reqURL, err := c.endpointURL(path)
if err != nil {
return nil, err
}
var bodyReader io.Reader
if len(body) > 0 {
bodyReader = bytes.NewReader(body)
}
req, err := http.NewRequestWithContext(ctx, method, reqURL, bodyReader)
if err != nil {
return nil, err
}
for k, vs := range headers {
for _, v := range vs {
req.Header.Add(k, v)
}
}
client := c.HTTP
if client == nil {
client = http.DefaultClient
}
if timeout > 0 {
ctx2, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
req = req.WithContext(ctx2)
}
resp, err := client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
raw, err := io.ReadAll(resp.Body)
if err != nil {
return nil, err
}
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return nil, fmt.Errorf("ilink %s %s: %d %s", method, path, resp.StatusCode, string(raw))
}
return raw, nil
}
// QRCodeResponse 获取二维码响应
type QRCodeResponse struct {
QRCode string `json:"qrcode"`
QRCodeImgContent string `json:"qrcode_img_content"`
}
// GetBotQRCode 获取绑定二维码
func (c *Client) GetBotQRCode(ctx context.Context, botType string, localTokenList []string) (*QRCodeResponse, error) {
if strings.TrimSpace(botType) == "" {
botType = DefaultBotType
}
body, _ := json.Marshal(map[string]interface{}{
"local_token_list": localTokenList,
})
path := "ilink/bot/get_bot_qrcode?bot_type=" + url.QueryEscape(botType)
raw, err := c.doRequest(ctx, http.MethodPost, path, body, c.authHeaders(), APIDefaultTimeout)
if err != nil {
return nil, err
}
var out QRCodeResponse
if err := json.Unmarshal(raw, &out); err != nil {
return nil, err
}
return &out, nil
}
// QRStatusResponse 二维码状态轮询响应
type QRStatusResponse struct {
Status string `json:"status"`
BotToken string `json:"bot_token"`
ILinkBotID string `json:"ilink_bot_id"`
ILinkUserID string `json:"ilink_user_id"`
BaseURL string `json:"baseurl"`
RedirectHost string `json:"redirect_host"`
}
// GetQRCodeStatus 长轮询二维码扫码状态
func (c *Client) GetQRCodeStatus(ctx context.Context, qrcode, verifyCode string) (*QRStatusResponse, error) {
path := "ilink/bot/get_qrcode_status?qrcode=" + url.QueryEscape(qrcode)
if verifyCode != "" {
path += "&verify_code=" + url.QueryEscape(verifyCode)
}
raw, err := c.doRequest(ctx, http.MethodGet, path, nil, c.commonHeaders(), QRLongPollTimeout)
if err != nil {
if ctx.Err() != nil {
return &QRStatusResponse{Status: "wait"}, nil
}
return &QRStatusResponse{Status: "wait"}, nil
}
var out QRStatusResponse
if err := json.Unmarshal(raw, &out); err != nil {
return nil, err
}
return &out, nil
}
// MessageItem 消息内容项
type MessageItem struct {
Type int `json:"type"`
TextItem *struct {
Text string `json:"text"`
} `json:"text_item,omitempty"`
}
// WeixinMessage 入站消息
type WeixinMessage struct {
FromUserID string `json:"from_user_id"`
MessageType int `json:"message_type"`
MessageState int `json:"message_state"`
ItemList []MessageItem `json:"item_list"`
ContextToken string `json:"context_token"`
}
// GetUpdatesResponse 长轮询消息响应
type GetUpdatesResponse struct {
Ret int `json:"ret"`
ErrCode int `json:"errcode"`
ErrMsg string `json:"errmsg"`
Msgs []WeixinMessage `json:"msgs"`
GetUpdatesBuf string `json:"get_updates_buf"`
LongPollingTimeoutMs int `json:"longpolling_timeout_ms"`
}
// GetUpdates 长轮询获取新消息
func (c *Client) GetUpdates(ctx context.Context, getUpdatesBuf string) (*GetUpdatesResponse, error) {
body, _ := json.Marshal(map[string]interface{}{
"get_updates_buf": getUpdatesBuf,
"base_info": c.buildBaseInfo(),
})
raw, err := c.doRequest(ctx, http.MethodPost, "ilink/bot/getupdates", body, c.authHeaders(), GetUpdatesTimeout)
if err != nil {
if ctx.Err() != nil {
return &GetUpdatesResponse{Ret: 0, GetUpdatesBuf: getUpdatesBuf}, nil
}
return &GetUpdatesResponse{Ret: 0, GetUpdatesBuf: getUpdatesBuf}, nil
}
var out GetUpdatesResponse
if err := json.Unmarshal(raw, &out); err != nil {
return nil, err
}
return &out, nil
}
// SendTextMessage 发送文本回复
func (c *Client) SendTextMessage(ctx context.Context, toUserID, contextToken, text, clientID string) error {
if clientID == "" {
clientID = randomClientID()
}
payload := map[string]interface{}{
"msg": map[string]interface{}{
"to_user_id": toUserID,
"client_id": clientID,
"message_type": 2,
"message_state": 2,
"context_token": contextToken,
"item_list": []map[string]interface{}{
{"type": 1, "text_item": map[string]string{"text": text}},
},
},
"base_info": c.buildBaseInfo(),
}
body, _ := json.Marshal(payload)
_, err := c.doRequest(ctx, http.MethodPost, "ilink/bot/sendmessage", body, c.authHeaders(), APIDefaultTimeout)
return err
}
func randomClientID() string {
var b [8]byte
_, _ = rand.Read(b[:])
return fmt.Sprintf("%x", b)
}
func sanitizeBotAgent(raw string) string {
raw = strings.TrimSpace(raw)
if raw == "" {
return DefaultBotAgent
}
if len(raw) > 256 {
return raw[:256]
}
return raw
}
// ExtractText 从消息中提取首条文本
func ExtractText(msg WeixinMessage) string {
for _, item := range msg.ItemList {
if item.Type == 1 && item.TextItem != nil {
return strings.TrimSpace(item.TextItem.Text)
}
}
return ""
}
-26
View File
@@ -1,26 +0,0 @@
package ilink
import (
"encoding/base64"
"fmt"
"strings"
"github.com/skip2/go-qrcode"
)
// QRCodeDataURL 将扫码内容(一般为 liteapp 链接)编码为 PNG data URL,供 Web 端展示。
// qrcode_img_content 不是图片直链,不能用作 <img src>。
func QRCodeDataURL(content string, size int) (string, error) {
content = strings.TrimSpace(content)
if content == "" {
return "", fmt.Errorf("empty qr content")
}
if size <= 0 {
size = 256
}
png, err := qrcode.Encode(content, qrcode.Medium, size)
if err != nil {
return "", err
}
return "data:image/png;base64," + base64.StdEncoding.EncodeToString(png), nil
}
-141
View File
@@ -1,141 +0,0 @@
package robot
import (
"context"
"encoding/json"
"strings"
"time"
"cyberstrike-ai/internal/config"
lark "github.com/larksuite/oapi-sdk-go/v3"
larkcore "github.com/larksuite/oapi-sdk-go/v3/core"
"github.com/larksuite/oapi-sdk-go/v3/event/dispatcher"
larkim "github.com/larksuite/oapi-sdk-go/v3/service/im/v1"
larkws "github.com/larksuite/oapi-sdk-go/v3/ws"
"go.uber.org/zap"
)
const (
larkReconnectInitial = 5 * time.Second // 首次重连间隔
larkReconnectMax = 60 * time.Second // 最大重连间隔
)
type larkTextContent struct {
Text string `json:"text"`
}
// StartLark 启动飞书长连接(无需公网),收到消息后调用 handler 并回复。
// 断线(如笔记本睡眠、网络中断)后会自动重连;ctx 被取消时退出,便于配置变更时重启。
func StartLark(ctx context.Context, robotsCfg config.RobotsConfig, h MessageHandler, logger *zap.Logger) {
cfg := robotsCfg.Lark
if !cfg.Enabled || cfg.AppID == "" || cfg.AppSecret == "" {
return
}
go runLarkLoop(ctx, cfg, robotsCfg.Session.StrictUserIdentityEnabled(), h, logger)
}
// runLarkLoop 循环维持飞书长连接:断开且 ctx 未取消时按退避间隔重连。
func runLarkLoop(ctx context.Context, cfg config.RobotLarkConfig, strictUserIdentity bool, h MessageHandler, logger *zap.Logger) {
backoff := larkReconnectInitial
for {
larkClient := lark.NewClient(cfg.AppID, cfg.AppSecret)
eventHandler := dispatcher.NewEventDispatcher("", "").OnP2MessageReceiveV1(func(ctx context.Context, event *larkim.P2MessageReceiveV1) error {
go handleLarkMessage(ctx, event, cfg, strictUserIdentity, h, larkClient, logger)
return nil
})
wsClient := larkws.NewClient(cfg.AppID, cfg.AppSecret,
larkws.WithEventHandler(eventHandler),
larkws.WithLogLevel(larkcore.LogLevelInfo),
)
logger.Info("飞书长连接正在连接…", zap.String("app_id", cfg.AppID))
err := wsClient.Start(ctx)
if ctx.Err() != nil {
logger.Info("飞书长连接已按配置重启关闭")
return
}
if err != nil {
logger.Warn("飞书长连接断开(如睡眠/断网),将自动重连", zap.Error(err), zap.Duration("retry_after", backoff))
}
select {
case <-ctx.Done():
return
case <-time.After(backoff):
if backoff < larkReconnectMax {
backoff *= 2
if backoff > larkReconnectMax {
backoff = larkReconnectMax
}
}
}
}
}
func handleLarkMessage(ctx context.Context, event *larkim.P2MessageReceiveV1, cfg config.RobotLarkConfig, strictUserIdentity bool, h MessageHandler, client *lark.Client, logger *zap.Logger) {
if event == nil || event.Event == nil || event.Event.Message == nil || event.Event.Sender == nil || event.Event.Sender.SenderId == nil {
return
}
msg := event.Event.Message
msgType := larkcore.StringValue(msg.MessageType)
if msgType != larkim.MsgTypeText {
logger.Debug("飞书暂仅处理文本消息", zap.String("msg_type", msgType))
return
}
var textBody larkTextContent
if err := json.Unmarshal([]byte(larkcore.StringValue(msg.Content)), &textBody); err != nil {
logger.Warn("飞书消息 Content 解析失败", zap.Error(err))
return
}
text := strings.TrimSpace(textBody.Text)
if text == "" {
return
}
userID := resolveLarkUserID(event, cfg.AllowChatIDFallback && !strictUserIdentity)
if userID == "" {
logger.Warn("飞书消息缺少可用用户标识,已忽略")
return
}
messageID := larkcore.StringValue(msg.MessageId)
reply := h.HandleMessage("lark", userID, text)
contentBytes, _ := json.Marshal(larkTextContent{Text: reply})
_, err := client.Im.Message.Reply(ctx, larkim.NewReplyMessageReqBuilder().
MessageId(messageID).
Body(larkim.NewReplyMessageReqBodyBuilder().
MsgType(larkim.MsgTypeText).
Content(string(contentBytes)).
Build()).
Build())
if err != nil {
logger.Warn("飞书回复失败", zap.String("message_id", messageID), zap.Error(err))
return
}
logger.Debug("飞书已回复", zap.String("message_id", messageID))
}
// resolveLarkUserID 提取飞书会话隔离键:
// tenant_key + 稳定用户标识(user_id/open_id/union_id);按配置可选 chat_id 兜底。
func resolveLarkUserID(event *larkim.P2MessageReceiveV1, allowChatIDFallback bool) string {
if event == nil || event.Event == nil || event.Event.Sender == nil || event.Event.Sender.SenderId == nil {
return ""
}
tenantKey := strings.TrimSpace(larkcore.StringValue(event.Event.Sender.TenantKey))
if tenantKey == "" {
tenantKey = "default"
}
prefix := "t:" + tenantKey + "|"
if id := strings.TrimSpace(larkcore.StringValue(event.Event.Sender.SenderId.UserId)); id != "" {
return prefix + "u:" + id
}
if id := strings.TrimSpace(larkcore.StringValue(event.Event.Sender.SenderId.OpenId)); id != "" {
return prefix + "o:" + id
}
if id := strings.TrimSpace(larkcore.StringValue(event.Event.Sender.SenderId.UnionId)); id != "" {
return prefix + "n:" + id
}
if allowChatIDFallback && event.Event.Message != nil {
if id := strings.TrimSpace(larkcore.StringValue(event.Event.Message.ChatId)); id != "" {
return prefix + "c:" + id
}
}
return ""
}
-192
View File
@@ -1,192 +0,0 @@
package robot
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"strconv"
"strings"
"time"
"cyberstrike-ai/internal/config"
"github.com/bwmarrin/discordgo"
lark "github.com/larksuite/oapi-sdk-go/v3"
larkim "github.com/larksuite/oapi-sdk-go/v3/service/im/v1"
"github.com/slack-go/slack"
)
// SendProactive sends a message without an inbound event. Platforms whose
// reply credentials are event-scoped deliberately return an error instead of
// pretending delivery succeeded.
func SendProactive(ctx context.Context, cfg config.RobotsConfig, platform, externalUserID, message string) error {
platform = strings.ToLower(strings.TrimSpace(platform))
userID := robotIdentityUserPart(externalUserID)
if userID == "" {
return fmt.Errorf("invalid robot recipient")
}
switch platform {
case "telegram":
if !cfg.Telegram.Enabled || strings.TrimSpace(cfg.Telegram.BotToken) == "" {
return fmt.Errorf("telegram is not configured")
}
id, err := strconv.ParseInt(userID, 10, 64)
if err != nil {
return fmt.Errorf("invalid telegram user id: %w", err)
}
return telegramSendReply(ctx, nilSafeHTTPClient(), strings.TrimSpace(cfg.Telegram.BotToken), id, message)
case "slack":
if !cfg.Slack.Enabled || strings.TrimSpace(cfg.Slack.BotToken) == "" {
return fmt.Errorf("slack is not configured")
}
api := slack.New(strings.TrimSpace(cfg.Slack.BotToken))
channel, _, _, err := api.OpenConversationContext(ctx, &slack.OpenConversationParameters{Users: []string{userID}})
if err != nil {
return err
}
for _, chunk := range splitTextChunks(message, slackMaxMessageRunes) {
if _, _, err = api.PostMessageContext(ctx, channel.ID, slack.MsgOptionText(chunk, false)); err != nil {
return err
}
}
return nil
case "discord":
if !cfg.Discord.Enabled || strings.TrimSpace(cfg.Discord.BotToken) == "" {
return fmt.Errorf("discord is not configured")
}
token := strings.TrimSpace(cfg.Discord.BotToken)
if !strings.HasPrefix(token, "Bot ") {
token = "Bot " + token
}
session, err := discordgo.New(token)
if err != nil {
return err
}
channel, err := session.UserChannelCreate(userID)
if err != nil {
return err
}
for _, chunk := range splitTextChunks(message, discordMaxMessageRunes) {
if _, err = session.ChannelMessageSend(channel.ID, chunk); err != nil {
return err
}
}
return nil
case "wecom":
return sendWecomProactive(ctx, cfg.Wecom, userID, message)
case "lark":
return sendLarkProactive(ctx, cfg.Lark, externalUserID, message)
default:
return fmt.Errorf("platform %s does not support proactive alerts yet", platform)
}
}
func SupportsProactive(platform string) bool {
switch strings.ToLower(strings.TrimSpace(platform)) {
case "telegram", "slack", "discord", "wecom", "lark":
return true
default:
return false
}
}
func nilSafeHTTPClient() *http.Client { return &http.Client{Timeout: 15 * time.Second} }
func sendWecomProactive(ctx context.Context, cfg config.RobotWecomConfig, userID, message string) error {
if !cfg.Enabled || strings.TrimSpace(cfg.CorpID) == "" || strings.TrimSpace(cfg.Secret) == "" || cfg.AgentID == 0 {
return fmt.Errorf("wecom proactive API is not configured")
}
tokenURL := "https://qyapi.weixin.qq.com/cgi-bin/gettoken?corpid=" + cfg.CorpID + "&corpsecret=" + cfg.Secret
req, err := http.NewRequestWithContext(ctx, http.MethodGet, tokenURL, nil)
if err != nil {
return err
}
resp, err := nilSafeHTTPClient().Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
body, err := io.ReadAll(resp.Body)
if err != nil {
return err
}
var tokenResp struct {
ErrCode int `json:"errcode"`
ErrMsg string `json:"errmsg"`
AccessToken string `json:"access_token"`
}
if err := json.Unmarshal(body, &tokenResp); err != nil {
return err
}
if tokenResp.ErrCode != 0 || tokenResp.AccessToken == "" {
return fmt.Errorf("wecom token: %s", tokenResp.ErrMsg)
}
payload, _ := json.Marshal(map[string]interface{}{
"touser": userID, "msgtype": "text", "agentid": cfg.AgentID,
"text": map[string]string{"content": message}, "safe": 0,
})
sendReq, err := http.NewRequestWithContext(ctx, http.MethodPost, "https://qyapi.weixin.qq.com/cgi-bin/message/send?access_token="+tokenResp.AccessToken, bytes.NewReader(payload))
if err != nil {
return err
}
sendReq.Header.Set("Content-Type", "application/json")
sendResp, err := nilSafeHTTPClient().Do(sendReq)
if err != nil {
return err
}
defer sendResp.Body.Close()
result, _ := io.ReadAll(sendResp.Body)
var parsed struct {
ErrCode int `json:"errcode"`
ErrMsg string `json:"errmsg"`
}
if err := json.Unmarshal(result, &parsed); err != nil {
return err
}
if parsed.ErrCode != 0 {
return fmt.Errorf("wecom send: %s", parsed.ErrMsg)
}
return nil
}
func sendLarkProactive(ctx context.Context, cfg config.RobotLarkConfig, identity, message string) error {
if !cfg.Enabled || strings.TrimSpace(cfg.AppID) == "" || strings.TrimSpace(cfg.AppSecret) == "" {
return fmt.Errorf("lark is not configured")
}
receiveIDType, receiveID := "user_id", robotIdentityUserPart(identity)
if idx := strings.LastIndex(identity, "|o:"); idx >= 0 {
receiveIDType, receiveID = "open_id", strings.TrimSpace(identity[idx+3:])
}
if idx := strings.LastIndex(identity, "|n:"); idx >= 0 {
receiveIDType, receiveID = "union_id", strings.TrimSpace(identity[idx+3:])
}
if receiveID == "" {
return fmt.Errorf("invalid lark recipient")
}
content, _ := json.Marshal(larkTextContent{Text: message})
client := lark.NewClient(cfg.AppID, cfg.AppSecret)
resp, err := client.Im.Message.Create(ctx, larkim.NewCreateMessageReqBuilder().
ReceiveIdType(receiveIDType).
Body(larkim.NewCreateMessageReqBodyBuilder().ReceiveId(receiveID).MsgType(larkim.MsgTypeText).Content(string(content)).Build()).Build())
if err != nil {
return err
}
if resp == nil || !resp.Success() {
return fmt.Errorf("lark send failed")
}
return nil
}
func robotIdentityUserPart(identity string) string {
identity = strings.TrimSpace(identity)
if i := strings.LastIndex(identity, "|u:"); i >= 0 {
return strings.TrimSpace(identity[i+3:])
}
if strings.HasPrefix(identity, "u:") {
return strings.TrimSpace(identity[2:])
}
return identity
}
-209
View File
@@ -1,209 +0,0 @@
package robot
import (
"context"
"strings"
"sync"
"cyberstrike-ai/internal/config"
"github.com/tencent-connect/botgo"
"github.com/tencent-connect/botgo/dto"
"github.com/tencent-connect/botgo/event"
"github.com/tencent-connect/botgo/openapi"
"github.com/tencent-connect/botgo/token"
"go.uber.org/zap"
)
const (
qqPlatform = "qq"
qqMaxMessageRunes = 3500
)
var (
qqHandlerMu sync.Mutex
qqHandler MessageHandler
qqLogger *zap.Logger
qqAPI openapi.OpenAPI
)
// StartQQ 启动 QQ 机器人 WebSocket(C2C 与群 @,出站连接,无需公网回调)。
func StartQQ(ctx context.Context, robotsCfg config.RobotsConfig, h MessageHandler, logger *zap.Logger) {
cfg := robotsCfg.QQ
if !cfg.Enabled || strings.TrimSpace(cfg.AppID) == "" || strings.TrimSpace(cfg.ClientSecret) == "" {
return
}
go runQQLoop(ctx, cfg, h, logger)
}
func runQQLoop(ctx context.Context, cfg config.RobotQQConfig, h MessageHandler, logger *zap.Logger) {
backoff := reconnectInitial
for {
if ctx.Err() != nil {
logger.Info("QQ 机器人 WebSocket 已按配置关闭")
return
}
err := runQQSession(ctx, cfg, h, logger)
if ctx.Err() != nil {
return
}
if err != nil {
logger.Warn("QQ 机器人 WebSocket 异常,将自动重连", zap.Error(err), zap.Duration("retry_after", backoff))
}
if !waitReconnect(ctx, &backoff) {
return
}
}
}
func runQQSession(ctx context.Context, cfg config.RobotQQConfig, h MessageHandler, logger *zap.Logger) error {
appID := strings.TrimSpace(cfg.AppID)
secret := strings.TrimSpace(cfg.ClientSecret)
credentials := &token.QQBotCredentials{AppID: appID, AppSecret: secret}
tokenSource := token.NewQQBotTokenSource(credentials)
if err := token.StartRefreshAccessToken(ctx, tokenSource); err != nil {
return err
}
var api openapi.OpenAPI
if cfg.Sandbox {
api = botgo.NewSandboxOpenAPI(appID, tokenSource)
} else {
api = botgo.NewOpenAPI(appID, tokenSource)
}
qqHandlerMu.Lock()
qqHandler = h
qqLogger = logger
qqAPI = api
qqHandlerMu.Unlock()
defer func() {
qqHandlerMu.Lock()
qqHandler = nil
qqLogger = nil
qqAPI = nil
qqHandlerMu.Unlock()
}()
intents := event.RegisterHandlers(
event.C2CMessageEventHandler(handleQQC2CMessage),
event.GroupATMessageEventHandler(handleQQGroupATMessage),
)
wsInfo, err := api.WS(ctx, nil, "")
if err != nil {
return err
}
logger.Info("QQ 机器人 WebSocket 正在连接…", zap.String("app_id", appID), zap.Bool("sandbox", cfg.Sandbox))
done := make(chan error, 1)
go func() {
done <- botgo.NewSessionManager().Start(wsInfo, tokenSource, &intents)
}()
select {
case <-ctx.Done():
return ctx.Err()
case err := <-done:
if err != nil {
return err
}
return nil
}
}
func handleQQC2CMessage(payload *dto.WSPayload, data *dto.WSC2CMessageData) error {
if data == nil || data.Author == nil {
return nil
}
text := strings.TrimSpace(data.Content)
if text == "" {
return nil
}
userOpenID := strings.TrimSpace(data.Author.ID)
if userOpenID == "" {
return nil
}
userID := "u:" + userOpenID
qqHandlerMu.Lock()
h := qqHandler
logger := qqLogger
api := qqAPI
qqHandlerMu.Unlock()
if h == nil || api == nil {
return nil
}
logger.Info("QQ 收到 C2C 消息", zap.String("from", userID), zap.String("content", text))
reply := h.HandleMessage(qqPlatform, userID, text)
return qqPostC2CReply(context.Background(), api, userOpenID, payload, data.ID, reply, logger)
}
func handleQQGroupATMessage(payload *dto.WSPayload, data *dto.WSGroupATMessageData) error {
if data == nil || data.Author == nil {
return nil
}
text := strings.TrimSpace(data.Content)
if text == "" {
return nil
}
userOpenID := strings.TrimSpace(data.Author.ID)
groupID := strings.TrimSpace(data.GroupID)
if userOpenID == "" {
return nil
}
userID := "g:" + groupID + "|u:" + userOpenID
qqHandlerMu.Lock()
h := qqHandler
logger := qqLogger
api := qqAPI
qqHandlerMu.Unlock()
if h == nil || api == nil {
return nil
}
logger.Info("QQ 收到群 @ 消息", zap.String("from", userID), zap.String("content", text))
reply := h.HandleMessage(qqPlatform, userID, text)
return qqPostGroupReply(context.Background(), api, groupID, payload, data.ID, reply, logger)
}
func qqPostC2CReply(ctx context.Context, api openapi.OpenAPI, userOpenID string, payload *dto.WSPayload, msgID, reply string, logger *zap.Logger) error {
reply = trimReply(reply)
if reply == "" {
return nil
}
for _, chunk := range splitTextChunks(reply, qqMaxMessageRunes) {
msg := &dto.MessageToCreate{
Content: chunk,
MsgID: msgID,
}
if payload != nil && payload.EventID != "" {
msg.EventID = payload.EventID
}
if _, err := api.PostC2CMessage(ctx, userOpenID, msg); err != nil {
logger.Warn("QQ 发送 C2C 回复失败", zap.String("to", userOpenID), zap.Error(err))
return err
}
}
return nil
}
func qqPostGroupReply(ctx context.Context, api openapi.OpenAPI, groupID string, payload *dto.WSPayload, msgID, reply string, logger *zap.Logger) error {
reply = trimReply(reply)
if reply == "" {
return nil
}
for _, chunk := range splitTextChunks(reply, qqMaxMessageRunes) {
msg := &dto.MessageToCreate{
Content: chunk,
MsgID: msgID,
}
if payload != nil && payload.EventID != "" {
msg.EventID = payload.EventID
}
if _, err := api.PostGroupMessage(ctx, groupID, msg); err != nil {
logger.Warn("QQ 发送群消息回复失败", zap.String("group", groupID), zap.Error(err))
return err
}
}
return nil
}
-38
View File
@@ -1,38 +0,0 @@
package robot
import (
"context"
"time"
)
const (
reconnectInitial = 5 * time.Second
reconnectMax = 60 * time.Second
)
func waitReconnect(ctx context.Context, backoff *time.Duration) bool {
if ctx.Err() != nil {
return false
}
select {
case <-ctx.Done():
return false
case <-time.After(*backoff):
if *backoff < reconnectMax {
*backoff *= 2
if *backoff > reconnectMax {
*backoff = reconnectMax
}
}
return true
}
}
func bumpBackoff(backoff *time.Duration) {
if *backoff < reconnectMax {
*backoff *= 2
if *backoff > reconnectMax {
*backoff = reconnectMax
}
}
}
-135
View File
@@ -1,135 +0,0 @@
package robot
import (
"context"
"strings"
"cyberstrike-ai/internal/config"
"github.com/slack-go/slack"
"github.com/slack-go/slack/slackevents"
"github.com/slack-go/slack/socketmode"
"go.uber.org/zap"
)
const (
slackPlatform = "slack"
slackMaxMessageRunes = 3900
)
// StartSlack 启动 Slack Socket Mode(出站 WebSocket,无需公网回调)。
func StartSlack(ctx context.Context, robotsCfg config.RobotsConfig, h MessageHandler, logger *zap.Logger) {
cfg := robotsCfg.Slack
if !cfg.Enabled || strings.TrimSpace(cfg.BotToken) == "" || strings.TrimSpace(cfg.AppToken) == "" {
return
}
go runSlackLoop(ctx, cfg, h, logger)
}
func runSlackLoop(ctx context.Context, cfg config.RobotSlackConfig, h MessageHandler, logger *zap.Logger) {
backoff := reconnectInitial
for {
err := runSlackSocket(ctx, cfg, h, logger)
if ctx.Err() != nil {
logger.Info("Slack Socket Mode 已按配置关闭")
return
}
if err != nil {
logger.Warn("Slack Socket Mode 异常,将自动重连", zap.Error(err), zap.Duration("retry_after", backoff))
}
if !waitReconnect(ctx, &backoff) {
return
}
}
}
func runSlackSocket(ctx context.Context, cfg config.RobotSlackConfig, h MessageHandler, logger *zap.Logger) error {
api := slack.New(
strings.TrimSpace(cfg.BotToken),
slack.OptionAppLevelToken(strings.TrimSpace(cfg.AppToken)),
)
client := socketmode.New(api)
logger.Info("Slack Socket Mode 正在连接…")
go func() {
for evt := range client.Events {
switch evt.Type {
case socketmode.EventTypeEventsAPI:
eventsAPIEvent, ok := evt.Data.(slackevents.EventsAPIEvent)
if !ok {
continue
}
client.Ack(*evt.Request)
if eventsAPIEvent.Type != slackevents.CallbackEvent {
continue
}
switch ev := eventsAPIEvent.InnerEvent.Data.(type) {
case *slackevents.MessageEvent:
handleSlackMessage(ctx, api, eventsAPIEvent.TeamID, ev, h, logger)
case *slackevents.AppMentionEvent:
handleSlackAppMention(ctx, api, eventsAPIEvent.TeamID, ev, h, logger)
}
case socketmode.EventTypeConnecting:
logger.Info("Slack Socket Mode 正在连接…")
case socketmode.EventTypeConnected:
logger.Info("Slack Socket Mode 已连接,等待收消息")
}
}
}()
return client.RunContext(ctx)
}
func handleSlackMessage(ctx context.Context, api *slack.Client, teamID string, ev *slackevents.MessageEvent, h MessageHandler, logger *zap.Logger) {
if ev == nil || ev.BotID != "" || ev.SubType != "" {
return
}
if ev.ChannelType != "im" {
return
}
text := strings.TrimSpace(ev.Text)
if text == "" {
return
}
userID := slackSessionKey(teamID, ev.User)
logger.Info("Slack 收到消息", zap.String("from", userID), zap.String("content", text))
reply := h.HandleMessage(slackPlatform, userID, text)
slackPostReply(ctx, api, ev.Channel, reply, logger)
}
func handleSlackAppMention(ctx context.Context, api *slack.Client, teamID string, ev *slackevents.AppMentionEvent, h MessageHandler, logger *zap.Logger) {
if ev == nil || ev.BotID != "" {
return
}
text := strings.TrimSpace(ev.Text)
if text == "" {
return
}
userID := slackSessionKey(teamID, ev.User)
logger.Info("Slack 收到 @ 消息", zap.String("from", userID), zap.String("content", text))
reply := h.HandleMessage(slackPlatform, userID, text)
slackPostReply(ctx, api, ev.Channel, reply, logger)
}
func slackSessionKey(teamID, userID string) string {
teamID = strings.TrimSpace(teamID)
userID = strings.TrimSpace(userID)
if teamID == "" {
teamID = "default"
}
return "t:" + teamID + "|u:" + userID
}
func slackPostReply(ctx context.Context, api *slack.Client, channel, reply string, logger *zap.Logger) {
reply = trimReply(reply)
if reply == "" {
return
}
for _, chunk := range splitTextChunks(reply, slackMaxMessageRunes) {
_, _, err := api.PostMessageContext(ctx, channel, slack.MsgOptionText(chunk, false))
if err != nil {
logger.Warn("Slack 发送回复失败", zap.String("channel", channel), zap.Error(err))
return
}
}
}
-29
View File
@@ -1,29 +0,0 @@
package robot
import "strings"
// splitTextChunks splits text into chunks no longer than maxRunes (rune count).
func splitTextChunks(text string, maxRunes int) []string {
text = strings.TrimSpace(text)
if text == "" || maxRunes <= 0 {
return nil
}
runes := []rune(text)
if len(runes) <= maxRunes {
return []string{text}
}
var out []string
for len(runes) > 0 {
end := maxRunes
if end > len(runes) {
end = len(runes)
}
out = append(out, string(runes[:end]))
runes = runes[end:]
}
return out
}
func trimReply(s string) string {
return strings.TrimSpace(s)
}
-262
View File
@@ -1,262 +0,0 @@
package robot
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"strings"
"time"
"cyberstrike-ai/internal/config"
"go.uber.org/zap"
)
const (
telegramPlatform = "telegram"
telegramAPIBase = "https://api.telegram.org"
telegramLongPollSec = 30
telegramMaxMessageRunes = 4096
)
type telegramUpdate struct {
UpdateID int `json:"update_id"`
Message *telegramMessage `json:"message"`
}
type telegramMessage struct {
MessageID int64 `json:"message_id"`
Chat telegramChat `json:"chat"`
From *telegramUser `json:"from"`
Text string `json:"text"`
Entities []telegramEntity `json:"entities"`
}
type telegramChat struct {
ID int64 `json:"id"`
Type string `json:"type"`
Title string `json:"title"`
}
type telegramUser struct {
ID int64 `json:"id"`
Username string `json:"username"`
IsBot bool `json:"is_bot"`
}
type telegramEntity struct {
Type string `json:"type"`
Offset int `json:"offset"`
Length int `json:"length"`
}
type telegramGetUpdatesResp struct {
OK bool `json:"ok"`
Result []telegramUpdate `json:"result"`
Description string `json:"description"`
}
type telegramBotMe struct {
OK bool `json:"ok"`
Result struct {
ID int64 `json:"id"`
Username string `json:"username"`
} `json:"result"`
}
// StartTelegram 启动 Telegram Bot 长轮询(getUpdates,无需公网回调)。
func StartTelegram(ctx context.Context, robotsCfg config.RobotsConfig, h MessageHandler, logger *zap.Logger) {
cfg := robotsCfg.Telegram
if !cfg.Enabled || strings.TrimSpace(cfg.BotToken) == "" {
return
}
go runTelegramLoop(ctx, cfg, h, logger)
}
func runTelegramLoop(ctx context.Context, cfg config.RobotTelegramConfig, h MessageHandler, logger *zap.Logger) {
backoff := reconnectInitial
for {
err := runTelegramPoll(ctx, cfg, h, logger)
if ctx.Err() != nil {
logger.Info("Telegram 长轮询已按配置关闭")
return
}
if err != nil {
logger.Warn("Telegram 长轮询异常,将自动重连", zap.Error(err), zap.Duration("retry_after", backoff))
}
if !waitReconnect(ctx, &backoff) {
return
}
}
}
func runTelegramPoll(ctx context.Context, cfg config.RobotTelegramConfig, h MessageHandler, logger *zap.Logger) error {
token := strings.TrimSpace(cfg.BotToken)
botUsername := strings.TrimSpace(cfg.BotUsername)
if botUsername == "" {
if name, err := telegramGetMe(ctx, token); err != nil {
logger.Warn("Telegram getMe 失败", zap.Error(err))
} else {
botUsername = name
}
}
offset := cfg.UpdateOffset
logger.Info("Telegram 长轮询已启动", zap.String("bot", botUsername))
client := &http.Client{Timeout: telegramLongPollSec*time.Second + 10*time.Second}
for {
select {
case <-ctx.Done():
return ctx.Err()
default:
}
updates, err := telegramGetUpdates(ctx, client, token, offset)
if err != nil {
return err
}
for _, u := range updates {
next := int64(u.UpdateID) + 1
if next > offset {
offset = next
}
if u.Message == nil || u.Message.From == nil || u.Message.From.IsBot {
continue
}
text := strings.TrimSpace(u.Message.Text)
if text == "" {
continue
}
chatType := strings.ToLower(strings.TrimSpace(u.Message.Chat.Type))
if chatType != "private" {
if !cfg.AllowGroupMessages {
continue
}
if botUsername != "" && !telegramMentionsBot(text, u.Message.Entities, botUsername) {
continue
}
}
userID := telegramSessionKey(chatType, u.Message.Chat.ID, u.Message.From.ID)
logger.Info("Telegram 收到消息", zap.String("from", userID), zap.String("content", text))
reply := h.HandleMessage(telegramPlatform, userID, text)
if err := telegramSendReply(ctx, client, token, u.Message.Chat.ID, reply); err != nil {
logger.Warn("Telegram 发送回复失败", zap.String("to", userID), zap.Error(err))
}
}
}
}
func telegramSessionKey(chatType string, chatID, fromUserID int64) string {
if chatType == "private" {
return fmt.Sprintf("u:%d", fromUserID)
}
return fmt.Sprintf("g:%d|u:%d", chatID, fromUserID)
}
func telegramMentionsBot(text string, entities []telegramEntity, botUsername string) bool {
needle := "@" + strings.TrimPrefix(strings.ToLower(botUsername), "@")
lower := strings.ToLower(text)
if strings.Contains(lower, needle) {
return true
}
for _, e := range entities {
if e.Type != "mention" {
continue
}
if e.Offset < 0 || e.Length <= 0 || e.Offset+e.Length > len(text) {
continue
}
mention := strings.ToLower(text[e.Offset : e.Offset+e.Length])
if mention == needle {
return true
}
}
return false
}
func telegramAPIURL(token, method string) string {
return fmt.Sprintf("%s/bot%s/%s", telegramAPIBase, token, method)
}
func telegramGetMe(ctx context.Context, token string) (string, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, telegramAPIURL(token, "getMe"), nil)
if err != nil {
return "", err
}
resp, err := http.DefaultClient.Do(req)
if err != nil {
return "", err
}
defer resp.Body.Close()
body, _ := io.ReadAll(resp.Body)
var parsed telegramBotMe
if err := json.Unmarshal(body, &parsed); err != nil {
return "", err
}
if !parsed.OK {
return "", fmt.Errorf("getMe failed: %s", string(body))
}
return parsed.Result.Username, nil
}
func telegramGetUpdates(ctx context.Context, client *http.Client, token string, offset int64) ([]telegramUpdate, error) {
url := fmt.Sprintf("%s?timeout=%d&allowed_updates=%s", telegramAPIURL(token, "getUpdates"), telegramLongPollSec, `["message"]`)
if offset > 0 {
url += fmt.Sprintf("&offset=%d", offset)
}
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
if err != nil {
return nil, err
}
resp, err := client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
body, err := io.ReadAll(resp.Body)
if err != nil {
return nil, err
}
var parsed telegramGetUpdatesResp
if err := json.Unmarshal(body, &parsed); err != nil {
return nil, err
}
if !parsed.OK {
if parsed.Description != "" {
return nil, fmt.Errorf("getUpdates: %s", parsed.Description)
}
return nil, fmt.Errorf("getUpdates failed: %s", string(body))
}
return parsed.Result, nil
}
func telegramSendReply(ctx context.Context, client *http.Client, token string, chatID int64, reply string) error {
reply = trimReply(reply)
if reply == "" {
return nil
}
for _, chunk := range splitTextChunks(reply, telegramMaxMessageRunes) {
payload := map[string]interface{}{
"chat_id": chatID,
"text": chunk,
}
body, _ := json.Marshal(payload)
req, err := http.NewRequestWithContext(ctx, http.MethodPost, telegramAPIURL(token, "sendMessage"), bytes.NewReader(body))
if err != nil {
return err
}
req.Header.Set("Content-Type", "application/json")
resp, err := client.Do(req)
if err != nil {
return err
}
resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("sendMessage status %d", resp.StatusCode)
}
}
return nil
}
-96
View File
@@ -1,96 +0,0 @@
package robot
import (
"context"
"strings"
"time"
"cyberstrike-ai/internal/config"
"cyberstrike-ai/internal/robot/ilink"
"go.uber.org/zap"
)
const (
wechatReconnectInitial = 5 * time.Second
wechatReconnectMax = 60 * time.Second
wechatPlatform = "wechat"
)
// StartWechat 启动微信 iLink 长轮询(无需公网回调),收到消息后调用 handler 并回复。
func StartWechat(ctx context.Context, robotsCfg config.RobotsConfig, h MessageHandler, appVersion string, logger *zap.Logger) {
cfg := robotsCfg.Wechat
if !cfg.Enabled || cfg.BotToken == "" {
return
}
go runWechatLoop(ctx, cfg, h, appVersion, logger)
}
func runWechatLoop(ctx context.Context, cfg config.RobotWechatConfig, h MessageHandler, appVersion string, logger *zap.Logger) {
backoff := wechatReconnectInitial
for {
err := runWechatPoll(ctx, cfg, h, appVersion, logger)
if ctx.Err() != nil {
logger.Info("微信 iLink 长轮询已按配置关闭")
return
}
if err != nil {
logger.Warn("微信 iLink 长轮询异常,将自动重连", zap.Error(err), zap.Duration("retry_after", backoff))
}
select {
case <-ctx.Done():
return
case <-time.After(backoff):
if backoff < wechatReconnectMax {
backoff *= 2
if backoff > wechatReconnectMax {
backoff = wechatReconnectMax
}
}
}
}
}
func runWechatPoll(ctx context.Context, cfg config.RobotWechatConfig, h MessageHandler, appVersion string, logger *zap.Logger) error {
client := ilink.NewClient(cfg.BaseURL, cfg.BotToken, cfg.BotAgent, ilink.BuildClientVersion(appVersion))
buf := cfg.GetUpdatesBuf
logger.Info("微信 iLink 长轮询已启动", zap.String("ilink_bot_id", cfg.ILinkBotID))
for {
select {
case <-ctx.Done():
return ctx.Err()
default:
}
resp, err := client.GetUpdates(ctx, buf)
if err != nil {
return err
}
if resp.ErrCode != 0 && resp.Ret != 0 {
logger.Warn("微信 getUpdates 返回错误", zap.Int("errcode", resp.ErrCode), zap.String("errmsg", resp.ErrMsg))
}
if resp.GetUpdatesBuf != "" {
buf = resp.GetUpdatesBuf
}
for _, msg := range resp.Msgs {
if msg.MessageType != 1 {
continue
}
text := ilink.ExtractText(msg)
if text == "" {
continue
}
userID := strings.TrimSpace(msg.FromUserID)
if userID == "" {
continue
}
logger.Info("微信收到消息", zap.String("from", userID), zap.String("content", text))
reply := h.HandleMessage(wechatPlatform, userID, text)
if strings.TrimSpace(reply) == "" {
continue
}
if err := client.SendTextMessage(ctx, userID, msg.ContextToken, reply, ""); err != nil {
logger.Warn("微信发送回复失败", zap.String("to", userID), zap.Error(err))
}
}
}
}