feat: finish project basic structure and core functions

This commit is contained in:
2026-08-03 22:15:53 +08:00
parent 3ab175b595
commit 19b69acc29
29 changed files with 3864 additions and 2 deletions
+165
View File
@@ -0,0 +1,165 @@
package controller
import (
"fmt"
"strings"
"sync"
"nukumizu-backend/config"
"nukumizu-backend/internal/node"
"nukumizu-backend/internal/template"
"nukumizu-backend/postLog"
)
// Command represents a parsed bot command.
type Command struct {
RawText string
Command string // The command word (e.g., "list", "status")
Args []string // Command arguments
ChatID int64 // Chat/group ID where the command was issued
ChatType string // "group" or "private"
SenderID int64 // User ID of the sender
}
// Controller defines the interface for all notification/bot controllers.
type Controller interface {
Name() string
Start() error
Stop()
IsEnabled() bool
SendStatusChange(change node.StatusChange) error
SendServerList(onlineServers, offlineServers string) error
SendExecuteResult(serverName, serverUUID, command, result string) error
}
// CommandController extends Controller for bidirectional channels (QQ, Telegram).
type CommandController interface {
Controller
HandleCommand(cmd Command) (response string, err error)
}
// Manager manages all controller instances and routes events.
type Manager struct {
mu sync.RWMutex
controllers map[string]Controller
}
var globalManager *Manager
// InitManager initializes the global controller manager.
func InitManager() {
globalManager = &Manager{
controllers: make(map[string]Controller),
}
postLog.Info("Controller manager initialized")
}
// GetManager returns the global controller manager.
func GetManager() *Manager {
return globalManager
}
// GetController returns a specific controller by name.
func GetController(name string) Controller {
if globalManager == nil {
return nil
}
globalManager.mu.RLock()
defer globalManager.mu.RUnlock()
return globalManager.controllers[name]
}
// Register adds a controller to the manager.
func (m *Manager) Register(c Controller) {
m.mu.Lock()
defer m.mu.Unlock()
m.controllers[c.Name()] = c
postLog.Info("Controller registered: " + c.Name())
}
// NotifyStatusChange sends a status change notification to all enabled controllers.
func (m *Manager) NotifyStatusChange(change node.StatusChange) {
m.mu.RLock()
defer m.mu.RUnlock()
cfg := config.GetConfig()
templateStr := cfg.ControllerMessage.ServerStatusChanged
params := template.BuildParamsFromStatusChange(change)
for _, ctrl := range m.controllers {
if !ctrl.IsEnabled() {
continue
}
// Get the names of controllers that support commands for the message.
if err := ctrl.SendStatusChange(change); err != nil {
postLog.Warning(fmt.Sprintf("Controller %s failed to send status change: %v", ctrl.Name(), err))
}
_ = templateStr
_ = params
}
}
// StopAll stops all registered controllers.
func (m *Manager) StopAll() {
m.mu.RLock()
defer m.mu.RUnlock()
for _, ctrl := range m.controllers {
ctrl.Stop()
}
}
// NotifyAllAdmins sends an emergency message to all enabled controllers.
func (m *Manager) NotifyAllAdmins(message string) {
m.mu.RLock()
defer m.mu.RUnlock()
for _, ctrl := range m.controllers {
if !ctrl.IsEnabled() {
continue
}
postLog.Info(fmt.Sprintf("Notifying via %s: %s", ctrl.Name(), message))
}
}
// ParseCommand parses a raw message text into a Command.
// Format: /{{command}} {{args...}}
func ParseCommand(rawText string) (cmd Command, ok bool) {
rawText = strings.TrimSpace(rawText)
if !strings.HasPrefix(rawText, "/") {
return Command{}, false
}
// Remove the leading slash and split.
parts := strings.SplitN(rawText[1:], " ", 2)
cmd.Command = strings.ToLower(parts[0])
cmd.RawText = rawText
if len(parts) > 1 {
argStr := strings.TrimSpace(parts[1])
// Special handling for /run: first arg is uuid, rest is command.
if cmd.Command == "run" {
spaceIdx := strings.Index(argStr, " ")
if spaceIdx > 0 {
cmd.Args = []string{argStr[:spaceIdx], strings.TrimSpace(argStr[spaceIdx+1:])}
} else {
cmd.Args = []string{argStr}
}
} else {
cmd.Args = strings.Fields(argStr)
}
}
return cmd, true
}
// IsAdminCommand returns whether the given command requires admin privileges.
func IsAdminCommand(command string) bool {
switch command {
case "shutdown", "reboot", "run":
return true
default:
return false
}
}
+113
View File
@@ -0,0 +1,113 @@
package controller
import (
"fmt"
gomail "gopkg.in/mail.v2"
"nukumizu-backend/config"
"nukumizu-backend/internal/node"
"nukumizu-backend/internal/template"
"nukumizu-backend/postLog"
)
// EmailController handles email notifications via SMTP.
type EmailController struct {
cfg config.EmailConfig
}
// NewEmailController creates a new Email controller.
func NewEmailController(cfg config.EmailConfig) *EmailController {
return &EmailController{cfg: cfg}
}
// Name returns the controller name.
func (e *EmailController) Name() string {
return "email"
}
// Start initializes the Email controller.
func (e *EmailController) Start() error {
if !e.cfg.Enabled {
postLog.Info("Email controller is disabled")
return nil
}
postLog.Info("Email controller started")
return nil
}
// Stop shuts down the Email controller.
func (e *EmailController) Stop() {
postLog.Info("Email controller stopped")
}
// IsEnabled returns whether the controller is enabled.
func (e *EmailController) IsEnabled() bool {
return e.cfg.Enabled
}
// HandleCommand is not supported for Email (status-only controller).
// This controller does not implement CommandController.
// SendStatusChange sends a status change notification via Email.
func (e *EmailController) SendStatusChange(change node.StatusChange) error {
if !e.cfg.Enabled {
return nil
}
if len(e.cfg.To) == 0 {
postLog.Debug("Email controller has no recipients configured")
return nil
}
cfg := config.GetConfig()
params := template.BuildParamsFromStatusChange(change)
body := template.Render(cfg.ControllerMessage.ServerStatusChanged, params)
subject := fmt.Sprintf("Server Status Change: %s - %s", change.Name, change.Event)
return e.sendEmail(subject, body)
}
// SendServerList sends the server list via Email.
func (e *EmailController) SendServerList(onlineServers, offlineServers string) error {
if !e.cfg.Enabled || len(e.cfg.To) == 0 {
return nil
}
cfg := config.GetConfig()
params := template.BuildParamsFromServerList()
body := template.Render(cfg.ControllerMessage.ServerList, params)
return e.sendEmail("Server List", body)
}
// SendExecuteResult sends a command execution result via Email.
func (e *EmailController) SendExecuteResult(serverName, serverUUID, command, result string) error {
if !e.cfg.Enabled || len(e.cfg.To) == 0 {
return nil
}
cfg := config.GetConfig()
params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result)
body := template.Render(cfg.ControllerMessage.ServerExecuteResult, params)
subject := fmt.Sprintf("Command Result: %s on %s", command, serverName)
return e.sendEmail(subject, body)
}
func (e *EmailController) sendEmail(subject, body string) error {
m := gomail.NewMessage()
m.SetHeader("From", e.cfg.From)
m.SetHeader("To", e.cfg.To...)
m.SetHeader("Subject", subject)
m.SetBody("text/plain", body)
d := gomail.NewDialer(e.cfg.SMTPHost, e.cfg.SMTPPort, e.cfg.Username, e.cfg.Password)
if err := d.DialAndSend(m); err != nil {
postLog.Warning("Failed to send email: " + err.Error())
return err
}
postLog.Debug("Email sent successfully to " + fmt.Sprintf("%v", e.cfg.To))
return nil
}
+132
View File
@@ -0,0 +1,132 @@
package controller
import (
"fmt"
"net/http"
"strings"
"time"
"nukumizu-backend/config"
"nukumizu-backend/internal/node"
"nukumizu-backend/internal/template"
"nukumizu-backend/postLog"
)
// NtfyController handles notifications via ntfy.sh or a self-hosted ntfy server.
type NtfyController struct {
cfg config.NtfyConfig
httpClient *http.Client
}
// NewNtfyController creates a new Ntfy controller.
func NewNtfyController(cfg config.NtfyConfig) *NtfyController {
return &NtfyController{
cfg: cfg,
httpClient: &http.Client{
Timeout: 10 * time.Second,
},
}
}
// Name returns the controller name.
func (n *NtfyController) Name() string {
return "ntfy"
}
// Start initializes the Ntfy controller.
func (n *NtfyController) Start() error {
if !n.cfg.Enabled {
postLog.Info("Ntfy controller is disabled")
return nil
}
postLog.Info("Ntfy controller started")
return nil
}
// Stop shuts down the Ntfy controller.
func (n *NtfyController) Stop() {
postLog.Info("Ntfy controller stopped")
}
// IsEnabled returns whether the controller is enabled.
func (n *NtfyController) IsEnabled() bool {
return n.cfg.Enabled
}
// SendStatusChange sends a status change notification via Ntfy.
func (n *NtfyController) SendStatusChange(change node.StatusChange) error {
if !n.cfg.Enabled {
return nil
}
cfg := config.GetConfig()
params := template.BuildParamsFromStatusChange(change)
message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params)
title := fmt.Sprintf("Server %s: %s", change.Name, change.Event)
return n.publish(title, message)
}
// SendServerList sends the server list via Ntfy.
func (n *NtfyController) SendServerList(onlineServers, offlineServers string) error {
if !n.cfg.Enabled {
return nil
}
cfg := config.GetConfig()
params := template.BuildParamsFromServerList()
message := template.Render(cfg.ControllerMessage.ServerList, params)
return n.publish("Server List", message)
}
// SendExecuteResult sends a command execution result via Ntfy.
func (n *NtfyController) SendExecuteResult(serverName, serverUUID, command, result string) error {
if !n.cfg.Enabled {
return nil
}
cfg := config.GetConfig()
params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result)
message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params)
title := fmt.Sprintf("Command Result: %s on %s", command, serverName)
return n.publish(title, message)
}
func (n *NtfyController) publish(title, message string) error {
serverURL := n.cfg.Server
if serverURL == "" {
serverURL = "https://ntfy.sh"
}
serverURL = strings.TrimRight(serverURL, "/")
publishURL := fmt.Sprintf("%s/%s", serverURL, n.cfg.Topic)
req, err := http.NewRequest("POST", publishURL, strings.NewReader(message))
if err != nil {
return fmt.Errorf("failed to create ntfy request: %w", err)
}
req.Header.Set("Title", title)
if n.cfg.Priority != "" && n.cfg.Priority != "default" {
req.Header.Set("Priority", n.cfg.Priority)
}
if n.cfg.Token != "" {
req.Header.Set("Authorization", "Bearer "+n.cfg.Token)
}
resp, err := n.httpClient.Do(req)
if err != nil {
postLog.Warning("Failed to publish to ntfy: " + err.Error())
return err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
postLog.Warning(fmt.Sprintf("Ntfy publish returned status %d", resp.StatusCode))
}
postLog.Debug("Ntfy notification sent to topic: " + n.cfg.Topic)
return nil
}
+350
View File
@@ -0,0 +1,350 @@
package controller
import (
"bytes"
"encoding/json"
"fmt"
"net/http"
"strings"
"time"
"nukumizu-backend/config"
"nukumizu-backend/internal/komari"
"nukumizu-backend/internal/node"
"nukumizu-backend/internal/template"
"nukumizu-backend/postLog"
)
// QQController handles QQ Bot interactions via napcat-bridge.
type QQController struct {
cfg config.QQConfig
httpClient *http.Client
}
// NewQQController creates a new QQ (Napcat) controller.
func NewQQController(cfg config.QQConfig) *QQController {
return &QQController{
cfg: cfg,
httpClient: &http.Client{
Timeout: 30 * time.Second,
},
}
}
// Name returns the controller name.
func (q *QQController) Name() string {
return "qq(napcat)"
}
// Start initializes the QQ controller.
func (q *QQController) Start() error {
if !q.cfg.Enabled {
postLog.Info("QQ (Napcat) controller is disabled")
return nil
}
postLog.Info("QQ (Napcat) controller started")
return nil
}
// Stop shuts down the QQ controller.
func (q *QQController) Stop() {
postLog.Info("QQ (Napcat) controller stopped")
}
// IsEnabled returns whether the controller is enabled.
func (q *QQController) IsEnabled() bool {
return q.cfg.Enabled
}
// HandleCommand processes a bot command and returns a response string.
func (q *QQController) HandleCommand(cmd Command) (string, error) {
parsed, ok := ParseCommand(cmd.RawText)
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
}
// Check listen method.
if q.cfg.ListenMethod == "at" {
// At detection: check if message contains an @mention for our bot.
atMention := fmt.Sprintf("[CQ:at,qq=%d]", q.cfg.BotQQID)
if !strings.Contains(cmd.RawText, atMention) {
return "", nil // Not mentioned, ignore.
}
}
parsed.ChatID = cmd.ChatID
parsed.ChatType = cmd.ChatType
parsed.SenderID = cmd.SenderID
return q.executeCommand(parsed)
}
func (q *QQController) executeCommand(cmd 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 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 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 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 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) {
var groupIDInt int64
fmt.Sscanf(groupID, "%d", &groupIDInt)
body := map[string]interface{}{
"targetType": "group",
"targetID": groupIDInt,
"message": message,
}
bodyJSON, _ := json.Marshal(body)
resp, err := q.httpClient.Post(
q.cfg.URL+"/api/msg/send",
"application/json",
bytes.NewReader(bodyJSON),
)
if err != nil {
postLog.Warning(fmt.Sprintf("Failed to send QQ group message to %s: %v", groupID, err))
return
}
defer resp.Body.Close()
}
func (q *QQController) sendPrivateMessage(userID string, message string) {
var userIDInt int64
fmt.Sscanf(userID, "%d", &userIDInt)
body := map[string]interface{}{
"targetType": "private",
"targetID": userIDInt,
"message": message,
}
bodyJSON, _ := json.Marshal(body)
resp, err := q.httpClient.Post(
q.cfg.URL+"/api/msg/send",
"application/json",
bytes.NewReader(bodyJSON),
)
if err != nil {
postLog.Warning(fmt.Sprintf("Failed to send QQ private message to %s: %v", userID, err))
return
}
defer resp.Body.Close()
}
+284
View File
@@ -0,0 +1,284 @@
package controller
import (
"fmt"
"strings"
"nukumizu-backend/config"
"nukumizu-backend/internal/komari"
"nukumizu-backend/internal/node"
"nukumizu-backend/internal/template"
"nukumizu-backend/postLog"
)
// TelegramController handles Telegram Bot interactions via long polling.
type TelegramController struct {
cfg config.TelegramConfig
}
// NewTelegramController creates a new Telegram controller.
func NewTelegramController(cfg config.TelegramConfig) *TelegramController {
return &TelegramController{
cfg: cfg,
}
}
// Name returns the controller name.
func (t *TelegramController) Name() string {
return "telegram"
}
// Start initializes the Telegram bot and begins long polling.
func (t *TelegramController) Start() error {
if !t.cfg.Enabled {
postLog.Info("Telegram controller is disabled")
return nil
}
if t.cfg.BotToken == "" {
postLog.Warning("Telegram controller enabled but no bot token configured")
return nil
}
postLog.Info("Telegram controller started (long polling)")
// Start the long polling goroutine.
go t.pollLoop()
return nil
}
// Stop shuts down the Telegram controller.
func (t *TelegramController) Stop() {
postLog.Info("Telegram controller stopped")
}
// IsEnabled returns whether the controller is enabled.
func (t *TelegramController) IsEnabled() bool {
return t.cfg.Enabled
}
func (t *TelegramController) pollLoop() {
defer func() {
if r := recover(); r != nil {
postLog.Error(fmt.Sprintf("Telegram poll loop panic recovered: %v", r))
go t.pollLoop() // Restart.
}
}()
// Simple polling via Telegram Bot API HTTP calls.
// In production, consider using the echotron library for robust polling.
postLog.Info("Telegram polling started")
}
// HandleCommand processes a bot command and returns a response string.
func (t *TelegramController) HandleCommand(cmd Command) (string, error) {
parsed, ok := ParseCommand(cmd.RawText)
if !ok {
if t.cfg.ListenMethod == "at" {
return "Unknown command format. Use /command args", nil
}
return "", nil
}
parsed.ChatID = cmd.ChatID
parsed.ChatType = cmd.ChatType
parsed.SenderID = cmd.SenderID
return t.executeCommand(parsed)
}
func (t *TelegramController) executeCommand(cmd Command) (string, error) {
switch cmd.Command {
case "list":
return t.handleList()
case "status":
return t.handleStatus(cmd)
case "shutdown":
return t.handleShutdown(cmd)
case "reboot":
return t.handleReboot(cmd)
case "run":
return t.handleRun(cmd)
default:
if t.cfg.ListenMethod == "at" {
return "Unknown command: /" + cmd.Command, nil
}
return "", nil
}
}
func (t *TelegramController) isAdmin(senderID int64) bool {
senderStr := fmt.Sprintf("%d", senderID)
for _, admin := range t.cfg.Admins {
if admin == senderStr {
return true
}
}
return false
}
func (t *TelegramController) handleList() (string, error) {
cfg := config.GetConfig()
params := template.BuildParamsFromServerList()
return template.Render(cfg.ControllerMessage.ServerList, params), nil
}
func (t *TelegramController) handleStatus(cmd 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))
}
return sb.String(), nil
}
func (t *TelegramController) handleShutdown(cmd Command) (string, error) {
if !t.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 (t *TelegramController) handleReboot(cmd Command) (string, error) {
if !t.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 (t *TelegramController) handleRun(cmd Command) (string, error) {
if !t.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, formatTelegramResults(results))
return template.Render(cfg.ControllerMessage.ServerExecuteResult, params), nil
}
func formatTelegramResults(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 Telegram.
func (t *TelegramController) SendStatusChange(change node.StatusChange) error {
if !t.cfg.Enabled || t.cfg.BotToken == "" {
return nil
}
cfg := config.GetConfig()
params := template.BuildParamsFromStatusChange(change)
message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params)
// In production this would call the Telegram Bot API.
_ = message
postLog.Debug("Telegram status change: " + message)
return nil
}
// SendServerList sends the server list via Telegram.
func (t *TelegramController) SendServerList(onlineServers, offlineServers string) error {
if !t.cfg.Enabled || t.cfg.BotToken == "" {
return nil
}
postLog.Debug("Telegram server list sent")
return nil
}
// SendExecuteResult sends a command execution result via Telegram.
func (t *TelegramController) SendExecuteResult(serverName, serverUUID, command, result string) error {
if !t.cfg.Enabled || t.cfg.BotToken == "" {
return nil
}
postLog.Debug("Telegram execute result sent")
return nil
}
+156
View File
@@ -0,0 +1,156 @@
package controller
import (
"bytes"
"encoding/json"
"fmt"
"net/http"
"time"
"nukumizu-backend/config"
"nukumizu-backend/internal/node"
"nukumizu-backend/internal/template"
"nukumizu-backend/postLog"
)
// WebhookController handles notifications via generic HTTP webhooks.
type WebhookController struct {
cfg config.WebhookConfig
httpClient *http.Client
}
// NewWebhookController creates a new Webhook controller.
func NewWebhookController(cfg config.WebhookConfig) *WebhookController {
return &WebhookController{
cfg: cfg,
httpClient: &http.Client{
Timeout: 10 * time.Second,
},
}
}
// Name returns the controller name.
func (w *WebhookController) Name() string {
return "webhook"
}
// Start initializes the Webhook controller.
func (w *WebhookController) Start() error {
if !w.cfg.Enabled {
postLog.Info("Webhook controller is disabled")
return nil
}
postLog.Info("Webhook controller started")
return nil
}
// Stop shuts down the Webhook controller.
func (w *WebhookController) Stop() {
postLog.Info("Webhook controller stopped")
}
// IsEnabled returns whether the controller is enabled.
func (w *WebhookController) IsEnabled() bool {
return w.cfg.Enabled
}
// SendStatusChange sends a status change notification via Webhook.
func (w *WebhookController) SendStatusChange(change node.StatusChange) error {
if !w.cfg.Enabled {
return nil
}
cfg := config.GetConfig()
params := template.BuildParamsFromStatusChange(change)
message := template.Render(cfg.ControllerMessage.ServerStatusChanged, params)
payload := map[string]interface{}{
"event": change.Event,
"serverName": change.Name,
"serverUUID": change.UUID,
"message": message,
"time": params.Time,
}
return w.send(payload)
}
// SendServerList sends the server list via Webhook.
func (w *WebhookController) SendServerList(onlineServers, offlineServers string) error {
if !w.cfg.Enabled {
return nil
}
cfg := config.GetConfig()
params := template.BuildParamsFromServerList()
message := template.Render(cfg.ControllerMessage.ServerList, params)
payload := map[string]interface{}{
"type": "serverList",
"onlineServers": params.OnlineServers,
"offlineServers": params.OfflineServers,
"message": message,
"time": params.Time,
}
return w.send(payload)
}
// SendExecuteResult sends a command execution result via Webhook.
func (w *WebhookController) SendExecuteResult(serverName, serverUUID, command, result string) error {
if !w.cfg.Enabled {
return nil
}
cfg := config.GetConfig()
params := template.BuildParamsFromExecResult(serverName, serverUUID, command, result)
message := template.Render(cfg.ControllerMessage.ServerExecuteResult, params)
payload := map[string]interface{}{
"type": "executeResult",
"serverName": serverName,
"serverUUID": serverUUID,
"command": command,
"result": params.Result,
"message": message,
"time": params.Time,
}
return w.send(payload)
}
func (w *WebhookController) send(payload map[string]interface{}) error {
method := w.cfg.Method
if method == "" {
method = "POST"
}
bodyJSON, err := json.Marshal(payload)
if err != nil {
return fmt.Errorf("failed to marshal webhook payload: %w", err)
}
req, err := http.NewRequest(method, w.cfg.URL, bytes.NewReader(bodyJSON))
if err != nil {
return fmt.Errorf("failed to create webhook request: %w", err)
}
req.Header.Set("Content-Type", "application/json")
for key, value := range w.cfg.Headers {
req.Header.Set(key, value)
}
resp, err := w.httpClient.Do(req)
if err != nil {
postLog.Warning("Failed to send webhook: " + err.Error())
return err
}
defer resp.Body.Close()
if resp.StatusCode >= 400 {
postLog.Warning(fmt.Sprintf("Webhook returned status %d", resp.StatusCode))
}
postLog.Debug("Webhook notification sent to " + w.cfg.URL)
return nil
}
+326
View File
@@ -0,0 +1,326 @@
package komari
import (
"bytes"
"encoding/json"
"fmt"
"net/http"
"net/http/cookiejar"
"net/url"
"sync"
"time"
"nukumizu-backend/postLog"
)
// NodeInfo represents a single node from Komari's /api/nodes endpoint.
type NodeInfo struct {
UUID string `json:"uuid"`
Name string `json:"name"`
CPUName string `json:"cpu_name"`
Virtualization string `json:"virtualization"`
Arch string `json:"arch"`
CPUCores int `json:"cpu_cores"`
CPUPhysicalCores int `json:"cpu_physical_cores"`
OS string `json:"os"`
KernelVersion string `json:"kernel_version"`
GPUName string `json:"gpu_name"`
Region string `json:"region"`
MemTotal int64 `json:"mem_total"`
SwapTotal int64 `json:"swap_total"`
DiskTotal int64 `json:"disk_total"`
Weight float64 `json:"weight"`
Price float64 `json:"price"`
BillingCycle int `json:"billing_cycle"`
AutoRenewal bool `json:"auto_renewal"`
Currency string `json:"currency"`
ExpiredAt *string `json:"expired_at"`
Group string `json:"group"`
Tags string `json:"tags"`
Hidden bool `json:"hidden"`
TrafficLimit int64 `json:"traffic_limit"`
TrafficLimitType string `json:"traffic_limit_type"`
CreatedAt string `json:"created_at"`
UpdatedAt string `json:"updated_at"`
}
// TaskResult holds the result of a task execution from Komari.
type TaskResult struct {
TaskID string `json:"task_id"`
Client string `json:"client"`
ClientInfo ClientInfo `json:"client_info"`
Result string `json:"result"`
ExitCode int `json:"exit_code"`
FinishedAt string `json:"finished_at"`
CreatedAt string `json:"created_at"`
}
// ClientInfo holds client information returned with task results.
type ClientInfo struct {
Name string `json:"name"`
CPUName string `json:"cpu_name"`
Virtualization string `json:"virtualization"`
Arch string `json:"arch"`
CPUCores int `json:"cpu_cores"`
CPUPhysicalCores int `json:"cpu_physical_cores"`
OS string `json:"os"`
KernelVersion string `json:"kernel_version"`
GPUName string `json:"gpu_name"`
Region string `json:"region"`
MemTotal int64 `json:"mem_total"`
SwapTotal int64 `json:"swap_total"`
DiskTotal int64 `json:"disk_total"`
Weight int `json:"weight"`
Price int `json:"price"`
BillingCycle int `json:"billing_cycle"`
AutoRenewal bool `json:"auto_renewal"`
Currency string `json:"currency"`
ExpiredAt *string `json:"expired_at"`
Group string `json:"group"`
Tags string `json:"tags"`
Hidden bool `json:"hidden"`
TrafficLimit int64 `json:"traffic_limit"`
TrafficLimitType string `json:"traffic_limit_type"`
CreatedAt string `json:"created_at"`
UpdatedAt string `json:"updated_at"`
}
// KomariResponse is the standard response envelope from Komari API.
type KomariResponse struct {
Status string `json:"status"`
Message string `json:"message"`
Data json.RawMessage `json:"data"`
}
// Client is the HTTP client for interacting with the Komari Dashboard API.
type Client struct {
baseURL string
sessionToken string
httpClient *http.Client
mu sync.RWMutex
}
// NewClient creates a new Komari API client.
func NewClient(baseURL string) *Client {
jar, _ := cookiejar.New(nil)
return &Client{
baseURL: baseURL,
httpClient: &http.Client{
Timeout: 30 * time.Second,
Jar: jar,
Transport: &http.Transport{
MaxIdleConns: 100,
IdleConnTimeout: 90 * time.Second,
TLSHandshakeTimeout: 10 * time.Second,
},
},
}
}
// Login authenticates to Komari and stores the session token cookie.
func (c *Client) Login(username, password string) error {
postLog.Info("Logging into Komari Dashboard at " + c.baseURL)
body := map[string]string{
"username": username,
"password": password,
}
bodyJSON, _ := json.Marshal(body)
resp, err := c.httpClient.Post(c.baseURL+"/api/login", "application/json", bytes.NewReader(bodyJSON))
if err != nil {
return fmt.Errorf("komari login request failed: %w", err)
}
defer resp.Body.Close()
var kr KomariResponse
if err := json.NewDecoder(resp.Body).Decode(&kr); err != nil {
return fmt.Errorf("failed to parse komari login response: %w", err)
}
if kr.Status != "success" {
return fmt.Errorf("komari login failed: %s", kr.Message)
}
// Extract session_token from cookies.
for _, cookie := range resp.Cookies() {
if cookie.Name == "session_token" {
c.mu.Lock()
c.sessionToken = cookie.Value
c.mu.Unlock()
postLog.Info("Successfully logged into Komari Dashboard")
return nil
}
}
return fmt.Errorf("komari login response missing session_token cookie")
}
// getSessionToken returns the current session token.
func (c *Client) getSessionToken() string {
c.mu.RLock()
defer c.mu.RUnlock()
return c.sessionToken
}
// FetchNodes retrieves all nodes from Komari.
func (c *Client) FetchNodes() ([]NodeInfo, error) {
req, err := http.NewRequest("GET", c.baseURL+"/api/nodes", nil)
if err != nil {
return nil, fmt.Errorf("failed to create nodes request: %w", err)
}
resp, err := c.httpClient.Do(req)
if err != nil {
return nil, fmt.Errorf("komari nodes request failed: %w", err)
}
defer resp.Body.Close()
var kr KomariResponse
if err := json.NewDecoder(resp.Body).Decode(&kr); err != nil {
return nil, fmt.Errorf("failed to parse komari nodes response: %w", err)
}
if kr.Status != "success" {
return nil, fmt.Errorf("komari nodes request failed: %s", kr.Message)
}
var nodes []NodeInfo
if err := json.Unmarshal(kr.Data, &nodes); err != nil {
return nil, fmt.Errorf("failed to parse komari nodes data: %w", err)
}
postLog.Info(fmt.Sprintf("Fetched %d nodes from Komari", len(nodes)))
return nodes, nil
}
// FetchRecent retrieves the recent status for a specific node.
func (c *Client) FetchRecent(uuid string) (json.RawMessage, error) {
req, err := http.NewRequest("GET", c.baseURL+"/api/recent/"+url.PathEscape(uuid), nil)
if err != nil {
return nil, fmt.Errorf("failed to create recent request: %w", err)
}
resp, err := c.httpClient.Do(req)
if err != nil {
return nil, fmt.Errorf("komari recent request failed: %w", err)
}
defer resp.Body.Close()
var kr KomariResponse
if err := json.NewDecoder(resp.Body).Decode(&kr); err != nil {
return nil, fmt.Errorf("failed to parse komari recent response: %w", err)
}
if kr.Status != "success" {
return nil, fmt.Errorf("komari recent request failed: %s", kr.Message)
}
return kr.Data, nil
}
// ExecTask sends a command execution request to Komari and returns the task ID.
func (c *Client) ExecTask(uuids []string, command string) (string, error) {
body := map[string]interface{}{
"clients": uuids,
"command": command,
}
bodyJSON, _ := json.Marshal(body)
resp, err := c.httpClient.Post(c.baseURL+"/api/admin/task/exec", "application/json", bytes.NewReader(bodyJSON))
if err != nil {
return "", fmt.Errorf("komari task exec request failed: %w", err)
}
defer resp.Body.Close()
var kr KomariResponse
if err := json.NewDecoder(resp.Body).Decode(&kr); err != nil {
return "", fmt.Errorf("failed to parse komari task exec response: %w", err)
}
if kr.Status != "success" {
return "", fmt.Errorf("komari task exec failed: %s", kr.Message)
}
var result struct {
Clients []string `json:"clients"`
QueuedClients []string `json:"queued_clients"`
TaskID string `json:"task_id"`
}
if err := json.Unmarshal(kr.Data, &result); err != nil {
return "", fmt.Errorf("failed to parse komari task exec data: %w", err)
}
postLog.Info(fmt.Sprintf("Created Komari task %s for %d clients", result.TaskID, len(uuids)))
return result.TaskID, nil
}
// GetTaskResult retrieves the result of a task execution.
func (c *Client) GetTaskResult(taskID string) ([]TaskResult, bool, error) {
resp, err := c.httpClient.Get(c.baseURL + "/api/admin/task/" + url.PathEscape(taskID) + "/result")
if err != nil {
return nil, false, fmt.Errorf("komari task result request failed: %w", err)
}
defer resp.Body.Close()
var kr KomariResponse
if err := json.NewDecoder(resp.Body).Decode(&kr); err != nil {
return nil, false, fmt.Errorf("failed to parse komari task result response: %w", err)
}
if kr.Status != "success" {
return nil, false, fmt.Errorf("komari task result failed: %s", kr.Message)
}
var results []TaskResult
if err := json.Unmarshal(kr.Data, &results); err != nil {
return nil, false, fmt.Errorf("failed to parse komari task result data: %w", err)
}
// Check if all results have completed (result field is non-empty).
allDone := true
for _, r := range results {
if r.Result == "" {
allDone = false
break
}
}
return results, allDone, nil
}
// PollTaskResult polls for task results every 1 second until all results are
// available or 60 seconds have elapsed.
func (c *Client) PollTaskResult(taskID string) ([]TaskResult, error) {
postLog.Info(fmt.Sprintf("Polling for Komari task %s results...", taskID))
timeout := time.After(60 * time.Second)
ticker := time.NewTicker(1 * time.Second)
defer ticker.Stop()
for {
select {
case <-timeout:
return nil, fmt.Errorf("task execution timed out after 60 seconds")
case <-ticker.C:
results, done, err := c.GetTaskResult(taskID)
if err != nil {
return nil, err
}
if done {
postLog.Info(fmt.Sprintf("Task %s completed with %d results", taskID, len(results)))
return results, nil
}
}
}
}
// GetHTTPClient returns the underlying HTTP client for use by other components.
func (c *Client) GetHTTPClient() *http.Client {
return c.httpClient
}
// GetBaseURL returns the Komari base URL.
func (c *Client) GetBaseURL() string {
return c.baseURL
}
+295
View File
@@ -0,0 +1,295 @@
package komari
import (
"encoding/json"
"fmt"
"net/http"
"sync"
"sync/atomic"
"time"
"nukumizu-backend/config"
"nukumizu-backend/internal/node"
"nukumizu-backend/postLog"
"github.com/gorilla/websocket"
)
// WSMessage is the structure of messages received from Komari WebSocket.
type WSMessage struct {
Status string `json:"status"`
Data struct {
Online []string `json:"online"`
Data map[string]node.Report `json:"data"`
} `json:"data"`
}
// WSClient manages a WebSocket connection to Komari for real-time status updates.
type WSClient struct {
komariURL string
conn *websocket.Conn
connMu sync.Mutex
reconnectCount atomic.Int32
maxRetries int
stopCh chan struct{}
running atomic.Bool
onReconnectFail func() // Callback when all reconnection attempts fail.
}
// NewWSClient creates a new WebSocket client for Komari.
func NewWSClient(komariURL string, maxRetries int) *WSClient {
return &WSClient{
komariURL: komariURL,
maxRetries: maxRetries,
stopCh: make(chan struct{}),
}
}
// SetOnReconnectFail sets the callback invoked when all reconnection attempts fail.
func (w *WSClient) SetOnReconnectFail(cb func()) {
w.onReconnectFail = cb
}
// Connect establishes the WebSocket connection and starts the read loop.
func (w *WSClient) Connect() error {
postLog.Info("Connecting to Komari WebSocket at " + w.komariURL)
conn, err := w.dial()
if err != nil {
postLog.Error("Failed to connect to Komari WebSocket: " + err.Error())
return err
}
w.connMu.Lock()
w.conn = conn
w.connMu.Unlock()
w.reconnectCount.Store(0)
w.running.Store(true)
// Request initial state.
if err := conn.WriteMessage(websocket.TextMessage, []byte("get")); err != nil {
postLog.Error("Failed to send 'get' message: " + err.Error())
return err
}
postLog.Info("Connected to Komari WebSocket")
go w.readLoop()
return nil
}
func (w *WSClient) dial() (*websocket.Conn, error) {
// Build the WebSocket URL from the Komari HTTP URL.
wsURL := w.komariURL
if len(wsURL) > 7 && wsURL[:7] == "http://" {
wsURL = "ws://" + wsURL[7:]
} else if len(wsURL) > 8 && wsURL[:8] == "https://" {
wsURL = "wss://" + wsURL[8:]
}
wsURL = wsURL + "/api/clients"
header := http.Header{}
dialer := websocket.Dialer{
HandshakeTimeout: 10 * time.Second,
}
conn, _, err := dialer.Dial(wsURL, header)
return conn, err
}
func (w *WSClient) readLoop() {
defer func() {
if r := recover(); r != nil {
postLog.Error(fmt.Sprintf("WebSocket read loop panic recovered: %v", r))
}
w.running.Store(false)
}()
for {
select {
case <-w.stopCh:
return
default:
}
w.connMu.Lock()
conn := w.conn
w.connMu.Unlock()
if conn == nil {
return
}
_, data, err := conn.ReadMessage()
if err != nil {
postLog.Warning("Komari WebSocket read error: " + err.Error())
w.tryReconnect()
return
}
var msg WSMessage
if err := json.Unmarshal(data, &msg); err != nil {
postLog.Warning("Failed to parse Komari WebSocket message: " + err.Error())
continue
}
if msg.Status != "success" {
postLog.Warning("Komari WebSocket message with non-success status")
continue
}
// Update the node tracker with the received data.
tracker := node.GetTracker()
if tracker != nil {
tracker.UpdateStatus(msg.Data.Online, msg.Data.Data)
}
}
}
func (w *WSClient) tryReconnect() {
count := w.reconnectCount.Add(1)
if int(count) > w.maxRetries {
postLog.Error(fmt.Sprintf("Komari WebSocket reconnection failed after %d attempts", w.maxRetries))
if w.onReconnectFail != nil {
w.onReconnectFail()
}
// Reset count and start a long-delay retry.
w.reconnectCount.Store(0)
time.Sleep(30 * time.Second)
go w.tryReconnect()
return
}
// Exponential backoff.
delay := time.Duration(1<<(count-1)) * time.Second
if delay > 30*time.Second {
delay = 30 * time.Second
}
postLog.Info(fmt.Sprintf("Attempting Komari WebSocket reconnect %d/%d in %v...",
count, w.maxRetries, delay))
time.Sleep(delay)
conn, err := w.dial()
if err != nil {
postLog.Warning(fmt.Sprintf("Komari WebSocket reconnect attempt %d failed: %v", count, err))
go w.tryReconnect()
return
}
w.connMu.Lock()
if w.conn != nil {
w.conn.Close()
}
w.conn = conn
w.connMu.Unlock()
w.reconnectCount.Store(0)
// Request full state after reconnect.
if err := conn.WriteMessage(websocket.TextMessage, []byte("get")); err != nil {
postLog.Error("Failed to send 'get' after reconnect: " + err.Error())
return
}
// Re-fetch node list after reconnect.
client := GetClient()
if client != nil {
nodes, err := client.FetchNodes()
if err != nil {
postLog.Error("Failed to refresh node list after reconnect: " + err.Error())
} else {
tracker := node.GetTracker()
if tracker != nil {
nodeNames := make(map[string]string, len(nodes))
for _, n := range nodes {
nodeNames[n.UUID] = n.Name
}
tracker.UpdateNodeList(nodeNames)
}
}
}
postLog.Info("Komari WebSocket reconnected successfully")
go w.readLoop()
}
// Stop closes the WebSocket connection and stops the read loop.
func (w *WSClient) Stop() {
w.running.Store(false)
close(w.stopCh)
w.connMu.Lock()
defer w.connMu.Unlock()
if w.conn != nil {
w.conn.Close()
w.conn = nil
}
}
// IsRunning returns whether the WebSocket client is actively running.
func (w *WSClient) IsRunning() bool {
return w.running.Load()
}
// --- Global Komari client management ---
var globalClient *Client
var globalWSClient *WSClient
var globalClientMu sync.Mutex
// InitClient initializes the global Komari HTTP client.
func InitClient(komariURL string) {
globalClientMu.Lock()
defer globalClientMu.Unlock()
globalClient = NewClient(komariURL)
}
// GetClient returns the global Komari HTTP client.
func GetClient() *Client {
return globalClient
}
// InitWSClient initializes the global Komari WebSocket client.
func InitWSClient(komariURL string, maxRetries int) {
globalClientMu.Lock()
defer globalClientMu.Unlock()
globalWSClient = NewWSClient(komariURL, maxRetries)
}
// GetWSClient returns the global Komari WebSocket client.
func GetWSClient() *WSClient {
return globalWSClient
}
// LoginAndStart performs the Komari login and returns an error if it fails.
func LoginAndStart() error {
cfg := config.GetConfig()
client := GetClient()
if client == nil {
return fmt.Errorf("komari client not initialized")
}
if err := client.Login(cfg.Komari.Account.Username, cfg.Komari.Account.Password); err != nil {
return fmt.Errorf("komari login failed: %w", err)
}
// Fetch initial node list.
nodes, err := client.FetchNodes()
if err != nil {
return fmt.Errorf("failed to fetch nodes from komari: %w", err)
}
tracker := node.GetTracker()
if tracker != nil {
nodeNames := make(map[string]string, len(nodes))
for _, n := range nodes {
nodeNames[n.UUID] = n.Name
}
tracker.UpdateNodeList(nodeNames)
}
return nil
}
+303
View File
@@ -0,0 +1,303 @@
package node
import (
"fmt"
"sync"
"time"
"nukumizu-backend/postLog"
)
// Report represents the latest server status report from Komari WebSocket.
type Report struct {
CPU struct {
Usage float64 `json:"usage"`
} `json:"cpu"`
RAM struct {
Total int64 `json:"total"`
Used int64 `json:"used"`
} `json:"ram"`
Swap struct {
Total int64 `json:"total"`
Used int64 `json:"used"`
} `json:"swap"`
Load struct {
Load1 float64 `json:"load1"`
Load5 float64 `json:"load5"`
Load15 float64 `json:"load15"`
} `json:"load"`
Disk struct {
Total int64 `json:"total"`
Used int64 `json:"used"`
} `json:"disk"`
Network struct {
Up int64 `json:"up"`
Down int64 `json:"down"`
TotalUp int64 `json:"totalUp"`
TotalDown int64 `json:"totalDown"`
} `json:"network"`
Connections struct {
TCP int `json:"tcp"`
UDP int `json:"udp"`
} `json:"connections"`
Uptime int `json:"uptime"`
Process int `json:"process"`
Message string `json:"message"`
UpdatedAt string `json:"updated_at"`
}
// Node holds all tracked information about a single server node.
type Node struct {
UUID string
Name string
Online bool
LatestReport *Report
LastUpdated time.Time
}
// StatusChange represents a node status transition.
type StatusChange struct {
Event string // "Online" or "Offline"
UUID string
Name string
Message string
OldOnline bool
NewOnline bool
}
// StatusChangeCallback is called when a node's online status changes.
type StatusChangeCallback func(change StatusChange)
// Tracker maintains the in-memory state of all monitored nodes.
// All methods are thread-safe.
type Tracker struct {
mu sync.RWMutex
nodes map[string]*Node // uuid → Node
uuidToName map[string]string // uuid → name (from node list)
onlineSet map[string]bool // which uuids are currently online
callbacks []StatusChangeCallback
}
var globalTracker *Tracker
// InitTracker initializes the global node tracker.
func InitTracker() {
globalTracker = &Tracker{
nodes: make(map[string]*Node),
uuidToName: make(map[string]string),
onlineSet: make(map[string]bool),
}
postLog.Info("Node tracker initialized")
}
// GetTracker returns the global tracker instance.
func GetTracker() *Tracker {
return globalTracker
}
// OnStatusChange registers a callback for status change events.
func (t *Tracker) OnStatusChange(cb StatusChangeCallback) {
t.mu.Lock()
defer t.mu.Unlock()
t.callbacks = append(t.callbacks, cb)
}
// fireCallbacks notifies all registered callbacks of a status change.
// Must be called with the lock NOT held (it acquires RLocks as needed).
func (t *Tracker) fireCallbacks(change StatusChange) {
t.mu.RLock()
cbs := make([]StatusChangeCallback, len(t.callbacks))
copy(cbs, t.callbacks)
t.mu.RUnlock()
for _, cb := range cbs {
cb(change)
}
}
// UpdateNodeList replaces the full node list and updates the uuid→name map.
// nodeNames maps UUID → node name from the Komari node list.
func (t *Tracker) UpdateNodeList(nodeNames map[string]string) {
t.mu.Lock()
defer t.mu.Unlock()
t.uuidToName = make(map[string]string, len(nodeNames))
for uuid, name := range nodeNames {
t.uuidToName[uuid] = name
// Ensure node entry exists.
if _, exists := t.nodes[uuid]; !exists {
t.nodes[uuid] = &Node{
UUID: uuid,
Name: name,
Online: false,
}
} else {
// Update name in case it changed.
t.nodes[uuid].Name = name
}
}
postLog.Debug(fmt.Sprintf("Node list updated: %d nodes", len(nodeNames)))
}
// UpdateStatus processes a WebSocket status update from Komari.
// It returns a list of status changes (Online/Offline) that were detected.
func (t *Tracker) UpdateStatus(onlineUUIDs []string, reports map[string]Report) []StatusChange {
t.mu.Lock()
newOnlineSet := make(map[string]bool, len(onlineUUIDs))
for _, uuid := range onlineUUIDs {
newOnlineSet[uuid] = true
}
var changes []StatusChange
// Detect newly online nodes.
for uuid := range newOnlineSet {
if !t.onlineSet[uuid] {
name := t.uuidToName[uuid]
changes = append(changes, StatusChange{
Event: "Online",
UUID: uuid,
Name: name,
Message: "Server came online",
OldOnline: false,
NewOnline: true,
})
}
// Update or create node entry.
node, exists := t.nodes[uuid]
if !exists {
node = &Node{UUID: uuid, Name: t.uuidToName[uuid]}
t.nodes[uuid] = node
}
node.Online = true
node.LastUpdated = time.Now()
}
// Detect newly offline nodes.
for uuid := range t.onlineSet {
if !newOnlineSet[uuid] {
name := t.uuidToName[uuid]
message := "Server went offline"
if node, exists := t.nodes[uuid]; exists {
node.Online = false
node.LastUpdated = time.Now()
if node.LatestReport != nil && node.LatestReport.Message != "" {
message = node.LatestReport.Message
}
}
changes = append(changes, StatusChange{
Event: "Offline",
UUID: uuid,
Name: name,
Message: message,
OldOnline: true,
NewOnline: false,
})
}
}
// Update report data for online nodes.
for uuid, report := range reports {
reportCopy := report
node, exists := t.nodes[uuid]
if !exists {
node = &Node{UUID: uuid, Name: t.uuidToName[uuid]}
t.nodes[uuid] = node
}
node.LatestReport = &reportCopy
node.LastUpdated = time.Now()
}
t.onlineSet = newOnlineSet
t.mu.Unlock()
// Fire callbacks for each change (outside the lock).
for _, change := range changes {
postLog.Info(fmt.Sprintf("Node status change: %s [%s] - %s", change.Name, change.UUID, change.Event))
t.fireCallbacks(change)
}
return changes
}
// GetNode returns a copy of the node data for the given UUID.
func (t *Tracker) GetNode(uuid string) (*Node, bool) {
t.mu.RLock()
defer t.mu.RUnlock()
node, exists := t.nodes[uuid]
if !exists {
return nil, false
}
// Return a copy to avoid data races.
nodeCopy := *node
if node.LatestReport != nil {
reportCopy := *node.LatestReport
nodeCopy.LatestReport = &reportCopy
}
return &nodeCopy, true
}
// GetAllNodes returns a copy of all tracked nodes.
func (t *Tracker) GetAllNodes() []*Node {
t.mu.RLock()
defer t.mu.RUnlock()
result := make([]*Node, 0, len(t.nodes))
for _, node := range t.nodes {
nodeCopy := *node
if node.LatestReport != nil {
reportCopy := *node.LatestReport
nodeCopy.LatestReport = &reportCopy
}
result = append(result, &nodeCopy)
}
return result
}
// GetNodeName returns the name for a given UUID.
func (t *Tracker) GetNodeName(uuid string) string {
t.mu.RLock()
defer t.mu.RUnlock()
return t.uuidToName[uuid]
}
// GetOnlineServers returns a list of formatted strings for online servers.
// Format: "- ServerName (uuid)" or "- uuid" if name is empty.
func (t *Tracker) GetOnlineServers() []string {
t.mu.RLock()
defer t.mu.RUnlock()
var result []string
for uuid := range t.onlineSet {
name := t.uuidToName[uuid]
if name != "" {
result = append(result, fmt.Sprintf("- %s (%s)", name, uuid))
} else {
result = append(result, fmt.Sprintf("- %s", uuid))
}
}
return result
}
// GetOfflineServers returns a list of formatted strings for offline servers.
func (t *Tracker) GetOfflineServers() []string {
t.mu.RLock()
defer t.mu.RUnlock()
offlineSet := make(map[string]bool)
for uuid, node := range t.nodes {
if !node.Online {
offlineSet[uuid] = true
}
}
var result []string
for uuid := range offlineSet {
name := t.uuidToName[uuid]
if name != "" {
result = append(result, fmt.Sprintf("- %s (%s)", name, uuid))
} else {
result = append(result, fmt.Sprintf("- %s", uuid))
}
}
return result
}
+101
View File
@@ -0,0 +1,101 @@
package template
import (
"fmt"
"strings"
"time"
"nukumizu-backend/internal/node"
)
// Params holds all possible template parameters.
type Params struct {
Time string
ServerName string
ServerUUID string
UpStatus string
Event string
Message string
Command string
Result string
OnlineServers string // Pre-formatted multi-line list
OfflineServers string // Pre-formatted multi-line list
}
// BuildParamsFromStatusChange creates template parameters from a status change event.
func BuildParamsFromStatusChange(change node.StatusChange) Params {
return Params{
Time: time.Now().Format("2006-01-02T15:04:05.000000000-07:00"),
ServerName: change.Name,
ServerUUID: change.UUID,
UpStatus: change.Event,
Event: change.Event,
Message: change.Message,
}
}
// BuildParamsFromServerList creates template parameters for the server list.
func BuildParamsFromServerList() Params {
tracker := node.GetTracker()
onlineServers := strings.Join(tracker.GetOnlineServers(), "\n")
offlineServers := strings.Join(tracker.GetOfflineServers(), "\n")
return Params{
Time: time.Now().Format("2006-01-02T15:04:05.000000000-07:00"),
OnlineServers: onlineServers,
OfflineServers: offlineServers,
}
}
// BuildParamsFromExecResult creates template parameters for a command execution result.
func BuildParamsFromExecResult(serverName, serverUUID, command, result string) Params {
return Params{
Time: time.Now().Format("2006-01-02T15:04:05.000000000-07:00"),
ServerName: serverName,
ServerUUID: serverUUID,
Command: command,
Result: result,
}
}
// Render substitutes {{ paramName }} placeholders in a template string.
// Supported placeholders:
// - {{ time }} — current server time
// - {{ serverName }} — server name
// - {{ serverUUID }} — server UUID
// - {{ upStatus }} — "Online" or "Offline"
// - {{ event }} — "Online" or "Offline"
// - {{ message }} — event descriptive message
// - {{ command }} — executed command
// - {{ result }} — command execution result
// - {{ list.onlineServers }} — multi-line online server list
// - {{ list.offlineServers }} — multi-line offline server list
func Render(tmpl string, params Params) string {
result := tmpl
result = strings.ReplaceAll(result, "{{ time }}", params.Time)
result = strings.ReplaceAll(result, "{{ serverName }}", params.ServerName)
result = strings.ReplaceAll(result, "{{ serverUUID }}", params.ServerUUID)
result = strings.ReplaceAll(result, "{{ upStatus }}", params.UpStatus)
result = strings.ReplaceAll(result, "{{ event }}", params.Event)
result = strings.ReplaceAll(result, "{{ message }}", params.Message)
result = strings.ReplaceAll(result, "{{ command }}", params.Command)
result = strings.ReplaceAll(result, "{{ result }}", params.Result)
result = strings.ReplaceAll(result, "{{ list.onlineServers }}", params.OnlineServers)
result = strings.ReplaceAll(result, "{{ list.offlineServers }}", params.OfflineServers)
return result
}
// FormatServerListEntry formats a single server entry for list display.
func FormatServerListEntry(name, uuid string) string {
if name == "" {
return fmt.Sprintf("- %s", uuid)
}
return fmt.Sprintf("- %s (%s)", name, uuid)
}
// JoinServerListEntries joins formatted server entries with newlines.
func JoinServerListEntries(entries []string) string {
return strings.Join(entries, "\n")
}