chore(qq_napcat): rewrite all napcat_bridge logic into this project to replace adapt into napcat_bridge

This commit is contained in:
2026-08-05 17:07:29 +08:00
parent dc30675e0f
commit 1d09765a4a
9 changed files with 507 additions and 184 deletions
@@ -0,0 +1,315 @@
package qq_napcat
import (
"bytes"
"encoding/json"
"fmt"
"io"
"net/http"
"net/url"
"sync"
"time"
"nukumizu-backend/postLog"
"nukumizu-backend/config"
"github.com/gorilla/websocket"
)
// APIResponse mirrors NapCat's HTTP API response envelope.
type APIResponse struct {
Status string `json:"status"`
RetCode int `json:"retcode"`
Message string `json:"message"`
Data json.RawMessage `json:"data"`
}
// Client is the HTTP + WebSocket client for a single NapCat instance.
// It both listens for incoming OneBot events over WebSocket and issues
// outbound NapCat API calls over HTTP.
type Client struct {
addr string
port string
token string
httpClient *http.Client
connMu sync.Mutex
conn *websocket.Conn
stopCh chan struct{}
stopOnce sync.Once
}
// NewClient creates a NapCat client for the given host/port/token.
func NewClient(addr, port, token string) *Client {
return &Client{
addr: addr,
port: port,
token: token,
httpClient: &http.Client{
Timeout: 30 * time.Second,
},
stopCh: make(chan struct{}),
}
}
func (c *Client) baseURL() string {
return fmt.Sprintf("http://%s:%s", c.addr, c.port)
}
// sendRequest performs an HTTP call to NapCat, injecting access_token into POST
// bodies AND setting Authorization: Bearer (compatible with HTTP Server adapters
// created in the NapCat Web UI).
func (c *Client) sendRequest(method, endpoint string, body []byte) ([]byte, error) {
url := c.baseURL() + endpoint
// For POST requests, inject access_token into the JSON body before sending.
var finalBody []byte
if body != nil && c.token != "" {
var bodyMap map[string]interface{}
if err := json.Unmarshal(body, &bodyMap); err == nil {
bodyMap["access_token"] = c.token
finalBody, _ = json.Marshal(bodyMap)
}
}
if finalBody == nil {
finalBody = body
}
var req *http.Request
var err error
if finalBody != nil {
req, err = http.NewRequest(method, url, bytes.NewReader(finalBody))
} else {
req, err = http.NewRequest(method, url, nil)
}
if err != nil {
return nil, fmt.Errorf("failed to create request: %w", err)
}
req.Header.Set("Content-Type", "application/json")
if c.token != "" {
req.Header.Set("Authorization", "Bearer "+c.token)
}
resp, err := c.httpClient.Do(req)
if err != nil {
return nil, fmt.Errorf("failed to send request to napcat: %w", err)
}
defer resp.Body.Close()
respBody, err := io.ReadAll(resp.Body)
if err != nil {
return nil, fmt.Errorf("failed to read napcat response: %w", err)
}
return respBody, nil
}
// parseResponse unmarshals NapCat's raw response body into an APIResponse.
func (c *Client) parseResponse(respBody []byte, caller string) (*APIResponse, error) {
var napcatResp APIResponse
if err := json.Unmarshal(respBody, &napcatResp); err != nil {
return nil, fmt.Errorf("failed to unmarshal napcat response (%s): %w\n%s", caller, err, string(respBody))
}
return &napcatResp, nil
}
// SendMsg sends a message via NapCat.
// targetType: "group" or "private"
// targetID: QQ group ID or user ID
// msg: message content
// hasAt: whether to prepend an @mention
// atTargetID: the QQ ID to @mention
func (c *Client) SendMsg(targetType string, targetID int64, msg string, hasAt bool, atTargetID int64) (*APIResponse, error) {
// Build message with optional @mention.
message := msg
if hasAt && atTargetID > 0 {
message = fmt.Sprintf("[CQ:at,qq=%d] %s", atTargetID, msg)
}
var endpoint string
var napcatReq map[string]interface{}
switch targetType {
case "group":
endpoint = "/send_group_msg"
napcatReq = map[string]interface{}{
"group_id": targetID,
"message": message,
}
case "private":
endpoint = "/send_private_msg"
napcatReq = map[string]interface{}{
"user_id": targetID,
"message": message,
}
default:
return nil, fmt.Errorf("invalid targetType: %s, must be 'group' or 'private'", targetType)
}
reqBody, err := json.Marshal(napcatReq)
if err != nil {
return nil, fmt.Errorf("failed to marshal request: %w", err)
}
if config.GetConfig().Debug.ShowNapcatAction {
postLog.Debug(fmt.Sprintf("[Napcat] SendMsg -> %s (%s): %s", endpoint, targetType, message))
}
respBody, err := c.sendRequest(http.MethodPost, endpoint, reqBody)
if err != nil {
return nil, err
}
return c.parseResponse(respBody, "SendMsg")
}
// RecallMsg recalls (deletes) a sent message via NapCat.
func (c *Client) RecallMsg(msgID int64) (*APIResponse, error) {
reqBody, err := json.Marshal(map[string]interface{}{
"message_id": msgID,
})
if err != nil {
return nil, fmt.Errorf("failed to marshal request: %w", err)
}
if config.GetConfig().Debug.ShowNapcatAction {
postLog.Debug(fmt.Sprintf("[Napcat] RecallMsg -> /delete_msg: %d", msgID))
}
respBody, err := c.sendRequest(http.MethodPost, "/delete_msg", reqBody)
if err != nil {
return nil, err
}
return c.parseResponse(respBody, "RecallMsg")
}
// GetGroupList retrieves the list of joined groups from NapCat.
func (c *Client) GetGroupList() (*APIResponse, error) {
if config.GetConfig().Debug.ShowNapcatAction {
postLog.Debug("[Napcat] GetGroupList -> /get_group_list")
}
respBody, err := c.sendRequest(http.MethodGet, "/get_group_list", nil)
if err != nil {
return nil, err
}
return c.parseResponse(respBody, "GetGroupList")
}
// GetGroupInfo retrieves detailed information about a specific group from NapCat.
func (c *Client) GetGroupInfo(groupID int64) (*APIResponse, error) {
reqBody, err := json.Marshal(map[string]interface{}{
"group_id": groupID,
})
if err != nil {
return nil, fmt.Errorf("failed to marshal request: %w", err)
}
if config.GetConfig().Debug.ShowNapcatAction {
postLog.Debug(fmt.Sprintf("[Napcat] GetGroupInfo -> /get_group_info: %d", groupID))
}
respBody, err := c.sendRequest(http.MethodPost, "/get_group_info", reqBody)
if err != nil {
return nil, err
}
return c.parseResponse(respBody, "GetGroupInfo")
}
// GetFriendsList retrieves the friends list from NapCat.
func (c *Client) GetFriendsList() (*APIResponse, error) {
if config.GetConfig().Debug.ShowNapcatAction {
postLog.Debug("[Napcat] GetFriendsList -> /get_friend_list")
}
respBody, err := c.sendRequest(http.MethodGet, "/get_friend_list", nil)
if err != nil {
return nil, err
}
return c.parseResponse(respBody, "GetFriendsList")
}
// Listen connects to the NapCat WebSocket server and calls onEvent for every raw
// message. It blocks forever, reconnecting every 5s after a disconnect, until
// Stop is called. Run it in a goroutine.
func (c *Client) Listen(onEvent func(raw []byte)) {
for {
select {
case <-c.stopCh:
return
default:
}
c.listenOnce(onEvent)
postLog.Warning("[Napcat] WebSocket disconnected, reconnecting in 5s...")
select {
case <-c.stopCh:
return
case <-time.After(5 * time.Second):
}
}
}
// listenOnce connects to the NapCat WebSocket and reads events until the
// connection drops or Stop is called.
func (c *Client) listenOnce(onEvent func(raw []byte)) {
wsURL := fmt.Sprintf("ws://%s:%s/", c.addr, c.port)
header := http.Header{}
if c.token != "" {
wsURL += "?access_token=" + url.QueryEscape(c.token)
header.Set("Authorization", "Bearer "+c.token)
}
conn, _, err := websocket.DefaultDialer.Dial(wsURL, header)
if err != nil {
postLog.Error(fmt.Sprintf("[Napcat] Failed to connect to NapCat WebSocket: %v", err))
return
}
c.connMu.Lock()
c.conn = conn
c.connMu.Unlock()
defer func() {
c.connMu.Lock()
if c.conn == conn {
c.conn = nil
}
c.connMu.Unlock()
conn.Close()
}()
postLog.Info(fmt.Sprintf("[Napcat] Connected to NapCat WebSocket at %s", wsURL))
for {
select {
case <-c.stopCh:
return
default:
}
_, msgBytes, err := conn.ReadMessage()
if err != nil {
postLog.Error(fmt.Sprintf("[Napcat] WebSocket read error: %v", err))
return
}
onEvent(msgBytes)
}
}
// Stop closes any open WebSocket connection and unblocks the Listen reconnect loop.
func (c *Client) Stop() {
c.stopOnce.Do(func() { close(c.stopCh) })
c.connMu.Lock()
defer c.connMu.Unlock()
if c.conn != nil {
c.conn.Close()
c.conn = nil
}
}
+446
View File
@@ -0,0 +1,446 @@
package qq_napcat
import (
"encoding/json"
"fmt"
"strings"
"nukumizu-backend/config"
"nukumizu-backend/internal/controller"
"nukumizu-backend/internal/komari"
"nukumizu-backend/internal/node"
"nukumizu-backend/internal/template"
"nukumizu-backend/postLog"
)
// oneBotEvent mirrors a OneBot 11 event pushed over the NapCat WebSocket.
// Only "message" events are handled; notice/request/meta_event are ignored.
type oneBotEvent struct {
PostType string `json:"post_type"`
MessageType string `json:"message_type"`
GroupID int64 `json:"group_id"`
UserID int64 `json:"user_id"`
RawMessage string `json:"raw_message"`
Message json.RawMessage `json:"message"`
Sender struct {
UserID int64 `json:"user_id"`
Nickname string `json:"nickname"`
} `json:"sender"`
SelfID int64 `json:"self_id"`
SubType string `json:"sub_type"`
}
// QQController handles QQ Bot interactions by connecting directly to NapCat.
type QQController struct {
cfg config.QQConfig
napcatClient *Client
}
// NewQQController creates a new QQ (Napcat) controller.
func NewQQController(cfg config.QQConfig) *QQController {
q := &QQController{cfg: cfg}
if cfg.Enabled {
q.napcatClient = NewClient(cfg.NapcatAddr, cfg.NapcatPort, cfg.NapcatToken)
}
return q
}
// Name returns the controller name.
func (q *QQController) Name() string {
return "qq(napcat)"
}
// Start initializes the QQ controller and starts the NapCat WebSocket listener.
func (q *QQController) Start() error {
if !q.cfg.Enabled {
// postLog.Info("QQ (Napcat) controller is disabled")
return nil
}
if q.napcatClient == nil {
postLog.Warning("QQ (Napcat) controller enabled but NapCat client is nil")
return nil
}
postLog.Info("QQ (Napcat) controller started, connecting to NapCat WebSocket...")
go func() {
defer func() {
if r := recover(); r != nil {
postLog.Error(fmt.Sprintf("Napcat WS listener panic recovered: %v", r))
}
}()
q.napcatClient.Listen(q.handleNapcatEvent)
}()
return nil
}
// Stop shuts down the QQ controller and its WebSocket listener.
func (q *QQController) Stop() {
if q.napcatClient != nil {
q.napcatClient.Stop()
}
postLog.Info("QQ (Napcat) controller stopped")
}
// IsEnabled returns whether the controller is enabled.
func (q *QQController) IsEnabled() bool {
return q.cfg.Enabled
}
// handleNapcatEvent processes a raw OneBot event received from the NapCat WebSocket.
func (q *QQController) handleNapcatEvent(raw []byte) {
var ev oneBotEvent
if err := json.Unmarshal(raw, &ev); err != nil {
postLog.Debug("Failed to parse NapCat WS event: " + err.Error())
return
}
// Only handle message events; ignore notice/request/meta_event.
if ev.PostType != "message" {
return
}
// Ignore messages the bot itself sent (echo prevention).
if q.isSelfMessage(ev) {
return
}
// Determine chat information.
chatID := ev.GroupID
if ev.MessageType == "private" {
chatID = ev.UserID
}
// Build a command from the raw message.
cmd := controller.Command{
RawText: ev.RawMessage,
ChatID: chatID,
ChatType: ev.MessageType,
SenderID: ev.UserID,
}
response, err := q.HandleCommand(cmd)
if err != nil {
postLog.Error("Napcat command handling failed: " + err.Error())
return
}
if response == "" {
return
}
// Reply in the same chat.
if ev.MessageType == "private" {
if _, err := q.napcatClient.SendMsg("private", ev.UserID, response, false, 0); err != nil {
postLog.Warning("Failed to send NapCat private reply: " + err.Error())
}
} else {
if _, err := q.napcatClient.SendMsg("group", ev.GroupID, response, false, 0); err != nil {
postLog.Warning("Failed to send NapCat group reply: " + err.Error())
}
}
}
// isSelfMessage returns true when the event was generated by the bot itself.
func (q *QQController) isSelfMessage(ev oneBotEvent) bool {
if q.cfg.BotQQID > 0 && ev.SelfID == q.cfg.BotQQID {
return true
}
if ev.UserID == ev.SelfID {
return true
}
return false
}
// HandleCommand processes a bot command and returns a response string.
func (q *QQController) HandleCommand(cmd controller.Command) (string, error) {
cfg := config.GetConfig()
// Debug: when ShowNapcatMsg is enabled, echo the received message back in real-time.
if cfg.Debug.ShowNapcatMsg {
postLog.Debug("Napcat message received: " + cmd.RawText)
return cmd.RawText, nil
}
text := cmd.RawText
// Check listen method. In "at" mode the @mention must be stripped before
// parsing, otherwise raw_message like "[CQ:at,qq=123] /list" fails the
// "/" prefix check.
if q.cfg.ListenMethod == "at" {
atMention := fmt.Sprintf("[CQ:at,qq=%d]", q.cfg.BotQQID)
if !strings.Contains(text, atMention) {
return "", nil // Not mentioned, ignore.
}
text = strings.ReplaceAll(text, atMention, "")
}
parsed, ok := controller.ParseCommand(text)
if !ok {
// Not a command. In "global" mode, silently ignore.
// In "at" mode, this would be an error - but at detection happens at the message level.
return "", nil
}
parsed.ChatID = cmd.ChatID
parsed.ChatType = cmd.ChatType
parsed.SenderID = cmd.SenderID
return q.executeCommand(parsed)
}
func (q *QQController) executeCommand(cmd controller.Command) (string, error) {
// Validate against supported commands.
switch cmd.Command {
case "list":
return q.handleList()
case "status":
return q.handleStatus(cmd)
case "shutdown":
return q.handleShutdown(cmd)
case "reboot":
return q.handleReboot(cmd)
case "run":
return q.handleRun(cmd)
default:
if q.cfg.ListenMethod == "at" {
return "Unknown command: /" + cmd.Command, nil
}
return "", nil // Global mode: silently ignore unknown commands.
}
}
func (q *QQController) isAdmin(senderID int64) bool {
senderStr := fmt.Sprintf("%d", senderID)
for _, admin := range q.cfg.Admins {
if admin == senderStr {
return true
}
}
return false
}
func (q *QQController) handleList() (string, error) {
cfg := config.GetConfig()
params := template.BuildParamsFromServerList()
return template.Render(cfg.ControllerMessage.ServerList, params), nil
}
func (q *QQController) handleStatus(cmd controller.Command) (string, error) {
if len(cmd.Args) < 1 {
return "Usage: /status <uuid>", nil
}
uuid := cmd.Args[0]
tracker := node.GetTracker()
n, exists := tracker.GetNode(uuid)
if !exists {
return fmt.Sprintf("Server with UUID %s not found", uuid), nil
}
statusStr := "Offline"
if n.Online {
statusStr = "Online"
}
var sb strings.Builder
sb.WriteString(fmt.Sprintf("Server: %s (%s)\n", n.Name, n.UUID))
sb.WriteString(fmt.Sprintf("Status: %s\n", statusStr))
if n.LatestReport != nil {
r := n.LatestReport
sb.WriteString(fmt.Sprintf("CPU: %.2f%%\n", r.CPU.Usage))
sb.WriteString(fmt.Sprintf("RAM: %d / %d\n", r.RAM.Used, r.RAM.Total))
sb.WriteString(fmt.Sprintf("Disk: %d / %d\n", r.Disk.Used, r.Disk.Total))
sb.WriteString(fmt.Sprintf("Network: ↑%d ↓%d\n", r.Network.Up, r.Network.Down))
sb.WriteString(fmt.Sprintf("Uptime: %d seconds\n", r.Uptime))
sb.WriteString(fmt.Sprintf("Processes: %d\n", r.Process))
if r.Message != "" {
sb.WriteString(fmt.Sprintf("Message: %s\n", r.Message))
}
}
return sb.String(), nil
}
func (q *QQController) handleShutdown(cmd controller.Command) (string, error) {
if !q.isAdmin(cmd.SenderID) {
return "Permission denied: admin only", nil
}
if len(cmd.Args) < 1 {
return "Usage: /shutdown <uuid>", nil
}
uuid := cmd.Args[0]
client := komari.GetClient()
if client == nil {
return "Error: Komari client not initialized", nil
}
_, err := client.ExecTask([]string{uuid}, "shutdown")
if err != nil {
return fmt.Sprintf("Error: %v", err), nil
}
return fmt.Sprintf("Shutdown command sent to server %s", uuid), nil
}
func (q *QQController) handleReboot(cmd controller.Command) (string, error) {
if !q.isAdmin(cmd.SenderID) {
return "Permission denied: admin only", nil
}
if len(cmd.Args) < 1 {
return "Usage: /reboot <uuid>", nil
}
uuid := cmd.Args[0]
client := komari.GetClient()
if client == nil {
return "Error: Komari client not initialized", nil
}
_, err := client.ExecTask([]string{uuid}, "reboot")
if err != nil {
return fmt.Sprintf("Error: %v", err), nil
}
return fmt.Sprintf("Reboot command sent to server %s", uuid), nil
}
func (q *QQController) handleRun(cmd controller.Command) (string, error) {
if !q.isAdmin(cmd.SenderID) {
return "Permission denied: admin only", nil
}
if len(cmd.Args) < 2 {
return "Usage: /run <uuid|all> <command>", nil
}
uuidArg := cmd.Args[0]
command := cmd.Args[1]
client := komari.GetClient()
if client == nil {
return "Error: Komari client not initialized", nil
}
var uuids []string
if uuidArg == "all" {
tracker := node.GetTracker()
for _, n := range tracker.GetAllNodes() {
uuids = append(uuids, n.UUID)
}
} else {
uuids = []string{uuidArg}
}
taskID, err := client.ExecTask(uuids, command)
if err != nil {
return fmt.Sprintf("Error executing command: %v", err), nil
}
results, err := client.PollTaskResult(taskID)
if err != nil {
return fmt.Sprintf("Error getting results: %v", err), nil
}
cfg := config.GetConfig()
params := template.BuildParamsFromExecResult(uuidArg, uuidArg, command, formatTaskResults(results))
return template.Render(cfg.ControllerMessage.ServerExecuteResult, params), nil
}
func formatTaskResults(results []komari.TaskResult) string {
var sb strings.Builder
for _, r := range results {
sb.WriteString(fmt.Sprintf("--- %s ---\n", r.Client))
sb.WriteString(r.Result)
sb.WriteString(fmt.Sprintf("\nExit code: %d\n", r.ExitCode))
}
return sb.String()
}
// SendStatusChange sends a status change notification via QQ.
func (q *QQController) SendStatusChange(change node.StatusChange) error {
if !q.cfg.Enabled {
return nil
}
cfg := config.GetConfig()
params := template.BuildParamsFromStatusChange(change)
message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params)
// Send to trusted groups.
for _, groupID := range q.cfg.TrustedGroups {
q.sendGroupMessage(groupID, message)
}
// Send to admins via private message.
for _, adminID := range q.cfg.Admins {
q.sendPrivateMessage(adminID, message)
}
return nil
}
// SendServerList sends the server list via QQ.
func (q *QQController) SendServerList(onlineServers, offlineServers string) error {
if !q.cfg.Enabled {
return nil
}
cfg := config.GetConfig()
params := template.BuildParamsFromServerList()
message := template.Render(cfg.ControllerMessage.ServerList, params)
for _, groupID := range q.cfg.TrustedGroups {
q.sendGroupMessage(groupID, message)
}
return nil
}
// SendExecuteResult sends a command execution result via QQ.
func (q *QQController) SendExecuteResult(serverName, serverUUID, command, result string) error {
if !q.cfg.Enabled {
return nil
}
cfg := config.GetConfig()
params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result)
message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params)
for _, groupID := range q.cfg.TrustedGroups {
q.sendGroupMessage(groupID, message)
}
return nil
}
func (q *QQController) sendGroupMessage(groupID string, message string) {
if q.napcatClient == nil {
postLog.Warning("Cannot send QQ group message: NapCat client not initialized")
return
}
var groupIDInt int64
if _, err := fmt.Sscanf(groupID, "%d", &groupIDInt); err != nil || groupIDInt == 0 {
postLog.Warning(fmt.Sprintf("Invalid QQ group ID: %s", groupID))
return
}
if _, err := q.napcatClient.SendMsg("group", groupIDInt, message, false, 0); err != nil {
postLog.Warning(fmt.Sprintf("Failed to send QQ group message to %s: %v", groupID, err))
}
}
func (q *QQController) sendPrivateMessage(userID string, message string) {
if q.napcatClient == nil {
postLog.Warning("Cannot send QQ private message: NapCat client not initialized")
return
}
var userIDInt int64
if _, err := fmt.Sscanf(userID, "%d", &userIDInt); err != nil || userIDInt == 0 {
postLog.Warning(fmt.Sprintf("Invalid QQ user ID: %s", userID))
return
}
if _, err := q.napcatClient.SendMsg("private", userIDInt, message, false, 0); err != nil {
postLog.Warning(fmt.Sprintf("Failed to send QQ private message to %s: %v", userID, err))
}
}