diff --git a/.gitignore b/.gitignore index 64af4f8..0186ae4 100644 --- a/.gitignore +++ b/.gitignore @@ -1 +1,3 @@ -agent.md \ No newline at end of file +/.claude +agent.md +config.json \ No newline at end of file diff --git a/README.md b/README.md index 71cf7f4..bd4c901 100644 --- a/README.md +++ b/README.md @@ -1,2 +1,158 @@ -# chisato-server +# Nukumizu Backend +Remote server monitoring and command execution subsystem for [Komari](https://www.komari.wiki). + +## Overview + +Nukumizu connects to a Komari Dashboard instance to: +- Monitor server status in real-time via WebSocket +- Execute commands on remote servers via the Komari task API +- Provide Bot interfaces on QQ (via NapCat) and Telegram for interactive server control +- Send one-way status notifications via Email, Ntfy, and Webhook + +## Project Structure + +``` +nukumizu-backend/ +├── main.go # Entry point, startup sequence, graceful shutdown +├── router.go # HTTP route registration +├── config/ +│ └── config.go # Configuration loading, defaults +├── handler/ +│ ├── user.go # User login/register handlers +│ ├── server.go # Server list/status/exec handlers +│ ├── bot.go # Bot message receive handler +│ └── health.go # Health check endpoint +├── database/ +│ └── user.go # SQLite user database +├── utils/ +│ ├── auth.go # Token management, authentication +│ └── middleware.go # Rate limit, CORS, XSS protection +├── postLog/ # Logging subsystem +├── internal/ +│ ├── komari/ +│ │ ├── client.go # Komari HTTP API client +│ │ └── ws.go # Komari WebSocket client +│ ├── node/ +│ │ └── tracker.go # Thread-safe node state tracking +│ ├── controller/ +│ │ ├── controller.go # Controller interface & manager +│ │ ├── qq.go # QQ (Napcat) Bot controller +│ │ ├── telegram.go # Telegram Bot controller +│ │ ├── email.go # Email notification controller +│ │ ├── ntfy.go # Ntfy notification controller +│ │ └── webhook.go # Webhook notification controller +│ └── template/ +│ └── template.go # Message template engine +``` + +## Configuration + +Copy and modify `config.json` at the project root: + +```json +{ + "system": { + "debugMode": true, + "listenAddr": "0.0.0.0", + "listenPort": "8080" + }, + "komari": { + "dashboardURL": "http://127.0.0.1:25774", + "account": { + "username": "admin", + "password": "admin" + } + }, + "controllerMethod": { + "qq(napcat)": { + "enabled": false, + "url": "http://127.0.0.1:8081", + "token": "", + "botQQID": 0, + "listenMethod": "global", + "admins": [], + "trustedGroups": [] + }, + "telegram": { + "enabled": false, + "botToken": "", + "listenMethod": "global", + "admins": [], + "trustedGroups": [] + }, + "email": { + "enabled": false, + "smtpHost": "", + "smtpPort": 587, + "username": "", + "password": "", + "from": "", + "to": [], + "useTLS": true + }, + "ntfy": { + "enabled": false, + "server": "https://ntfy.sh", + "topic": "", + "token": "", + "priority": "default" + }, + "webhook": { + "enabled": false, + "url": "", + "method": "POST", + "headers": {}, + "template": "" + } + }, + "controllerMessage": { + "SERVER_STATUS_CHANGED": "Server Status Changed Alert\n{{ serverName }} - {{ upStatus }}\nEvent: {{ event }}\nServer Name: {{ serverName }}\nMessage: {{ message }}\nTime: {{ time }}", + "SERVER_LIST": "All server list:\nOnline:\n{{ list.onlineServers }}\nOffline:\n{{ list.offlineServers }}", + "SERVER_EXECUTE_RESULT": "Command execute result:\nServer Name: {{ serverName }}\nCommand: {{ command }}\n***Result***\n\n{{ result }}\n\n************\nTime: {{ time }}" + } +} +``` + +## API Endpoints + +All API responses follow the format: `{"success": bool, "message": "..."}` + +Authentication is via `X-Token` and `X-Timestamp` HTTP headers. + +| Endpoint | Method | Auth | Description | +|---|---|---|---| +| `/api/user/login` | POST | None | User login | +| `/api/user/register` | POST | None | First-time registration | +| `/api/server/list` | GET | bot/admin | List all servers | +| `/api/server/getStatus` | GET | bot/admin | Get server recent status | +| `/api/server/exec` | POST | bot/admin | Execute command on server(s) | +| `/api/bot/msg/recv` | POST | bot | Receive bot messages (from napcat-bridge) | +| `/health` | GET | None | Health check | +| `/api/system/getLogs` | WS | None | Real-time log streaming | + +## Bot Commands + +| Command | Permission | Description | +|---|---|---| +| `/list` | None | List all server status | +| `/shutdown ` | Admin | Shutdown specific server | +| `/reboot ` | Admin | Reboot specific server | +| `/status ` | None | Get server detailed status | +| `/run ` | Admin | Run command on server(s) | + +## Building + +```bash +go build -o nukumizu-backend . +``` + +## Running + +```bash +./nukumizu-backend -config config.json +``` + +## License + +See [LICENSE](LICENSE). diff --git a/config/config.go b/config/config.go new file mode 100644 index 0000000..2966aab --- /dev/null +++ b/config/config.go @@ -0,0 +1,110 @@ +package config + +import ( + "encoding/json" + "fmt" + "os" +) + +// LoadConfig reads and parses the configuration file, applies defaults, +// and stores it as a global singleton. +func LoadConfig(configPath string) (*Config, error) { + data, err := os.ReadFile(configPath) + if err != nil { + return nil, fmt.Errorf("failed to read config file: %w", err) + } + + var cfg Config + if err := json.Unmarshal(data, &cfg); err != nil { + return nil, fmt.Errorf("failed to parse config file: %w", err) + } + + // Apply defaults for System. + if cfg.System.ListenAddr == "" { + cfg.System.ListenAddr = "0.0.0.0" + } + if cfg.System.ListenPort == "" { + cfg.System.ListenPort = "8080" + } + + // Apply defaults for QQ controller. + if cfg.ControllerMethod.QQ.ListenMethod == "" { + cfg.ControllerMethod.QQ.ListenMethod = "global" + } + if cfg.ControllerMethod.QQ.Admins == nil { + cfg.ControllerMethod.QQ.Admins = []string{} + } + if cfg.ControllerMethod.QQ.TrustedGroups == nil { + cfg.ControllerMethod.QQ.TrustedGroups = []string{} + } + + // Apply defaults for Telegram controller. + if cfg.ControllerMethod.Telegram.ListenMethod == "" { + cfg.ControllerMethod.Telegram.ListenMethod = "global" + } + if cfg.ControllerMethod.Telegram.Admins == nil { + cfg.ControllerMethod.Telegram.Admins = []string{} + } + if cfg.ControllerMethod.Telegram.TrustedGroups == nil { + cfg.ControllerMethod.Telegram.TrustedGroups = []string{} + } + + // Apply defaults for Email controller. + if cfg.ControllerMethod.Email.SMTPPort == 0 { + cfg.ControllerMethod.Email.SMTPPort = 587 + } + if cfg.ControllerMethod.Email.To == nil { + cfg.ControllerMethod.Email.To = []string{} + } + + // Apply defaults for Ntfy controller. + if cfg.ControllerMethod.Ntfy.Server == "" { + cfg.ControllerMethod.Ntfy.Server = "https://ntfy.sh" + } + if cfg.ControllerMethod.Ntfy.Priority == "" { + cfg.ControllerMethod.Ntfy.Priority = "default" + } + + // Apply defaults for Webhook controller. + if cfg.ControllerMethod.Webhook.Method == "" { + cfg.ControllerMethod.Webhook.Method = "POST" + } + if cfg.ControllerMethod.Webhook.Headers == nil { + cfg.ControllerMethod.Webhook.Headers = map[string]string{} + } + + // Apply defaults for paths. + if cfg.DataPath == "" { + cfg.DataPath = "./data" + } + if cfg.DBPath == "" { + cfg.DBPath = "./logs.db" + } + + // Apply default message templates if not specified. + if cfg.ControllerMessage.ServerStatusChanged == "" { + cfg.ControllerMessage.ServerStatusChanged = "Server Status Changed Alert\n{{ serverName }} - {{ upStatus }}\nEvent: {{ event }}\nServer Name: {{ serverName }}\nMessage: {{ message }}\nTime: {{ time }}" + } + if cfg.ControllerMessage.ServerList == "" { + cfg.ControllerMessage.ServerList = "All server list:\nOnline:\n{{ list.onlineServers }}\nOffline:\n{{ list.offlineServers }}" + } + if cfg.ControllerMessage.ServerExecuteResult == "" { + cfg.ControllerMessage.ServerExecuteResult = "Command execute result:\nServer Name: {{ serverName }}\nCommand: {{ command }}\n***Result***\n\n{{ result }}\n\n************\nTime: {{ time }}" + } + + globalConfig = &cfg + return &cfg, nil +} + +// GetConfig returns the global configuration singleton. +func GetConfig() *Config { + return globalConfig +} + +// IsDebugMode returns whether debug mode is enabled. +func IsDebugMode() bool { + if globalConfig == nil { + return false + } + return globalConfig.System.DebugMode +} diff --git a/config/variables.go b/config/variables.go new file mode 100644 index 0000000..7cb049c --- /dev/null +++ b/config/variables.go @@ -0,0 +1,98 @@ +package config + +// SystemConfig holds system-level configuration. +type SystemConfig struct { + DebugMode bool `json:"debugMode"` + ListenAddr string `json:"listenAddr"` + ListenPort string `json:"listenPort"` +} + +// KomariAccount holds Komari login credentials. +type KomariAccount struct { + Username string `json:"username"` + Password string `json:"password"` +} + +// KomariConfig holds Komari Dashboard connection settings. +type KomariConfig struct { + DashboardURL string `json:"dashboardURL"` + Account KomariAccount `json:"account"` +} + +// QQConfig holds QQ (Napcat) Bot controller configuration. +type QQConfig struct { + Enabled bool `json:"enabled"` + URL string `json:"url"` + Token string `json:"token"` + BotQQID int64 `json:"botQQID"` + ListenMethod string `json:"listenMethod"` + Admins []string `json:"admins"` + TrustedGroups []string `json:"trustedGroups"` +} + +// TelegramConfig holds Telegram Bot controller configuration. +type TelegramConfig struct { + Enabled bool `json:"enabled"` + BotToken string `json:"botToken"` + ListenMethod string `json:"listenMethod"` + Admins []string `json:"admins"` + TrustedGroups []string `json:"trustedGroups"` +} + +// EmailConfig holds Email notification controller configuration. +type EmailConfig struct { + Enabled bool `json:"enabled"` + SMTPHost string `json:"smtpHost"` + SMTPPort int `json:"smtpPort"` + Username string `json:"username"` + Password string `json:"password"` + From string `json:"from"` + To []string `json:"to"` + UseTLS bool `json:"useTLS"` +} + +// NtfyConfig holds Ntfy notification controller configuration. +type NtfyConfig struct { + Enabled bool `json:"enabled"` + Server string `json:"server"` + Topic string `json:"topic"` + Token string `json:"token"` + Priority string `json:"priority"` +} + +// WebhookConfig holds Webhook notification controller configuration. +type WebhookConfig struct { + Enabled bool `json:"enabled"` + URL string `json:"url"` + Method string `json:"method"` + Headers map[string]string `json:"headers"` + Template string `json:"template"` +} + +// ControllerMethodConfig holds all controller method configurations. +type ControllerMethodConfig struct { + QQ QQConfig `json:"qq(napcat)"` + Telegram TelegramConfig `json:"telegram"` + Email EmailConfig `json:"email"` + Ntfy NtfyConfig `json:"ntfy"` + Webhook WebhookConfig `json:"webhook"` +} + +// ControllerMessageConfig holds message templates for controller responses. +type ControllerMessageConfig struct { + ServerStatusChanged string `json:"SERVER_STATUS_CHANGED"` + ServerList string `json:"SERVER_LIST"` + ServerExecuteResult string `json:"SERVER_EXECUTE_RESULT"` +} + +// Config is the top-level application configuration. +type Config struct { + System SystemConfig `json:"system"` + Komari KomariConfig `json:"komari"` + ControllerMethod ControllerMethodConfig `json:"controllerMethod"` + ControllerMessage ControllerMessageConfig `json:"controllerMessage"` + DataPath string `json:"dataPath"` + DBPath string `json:"dbPath"` +} + +var globalConfig *Config \ No newline at end of file diff --git a/database/user.go b/database/user.go new file mode 100644 index 0000000..663da77 --- /dev/null +++ b/database/user.go @@ -0,0 +1,140 @@ +package database + +import ( + "crypto/sha256" + "database/sql" + "encoding/hex" + "fmt" + "os" + "path/filepath" + "time" + + _ "modernc.org/sqlite" +) + +// UserDB is the global database handle for user data. +var UserDB *sql.DB + +// InitUserDB opens (or creates) the user database and initializes the users table. +func InitUserDB(dbPath string) error { + dir := filepath.Dir(dbPath) + if dir != "." { + if err := os.MkdirAll(dir, 0755); err != nil { + return fmt.Errorf("failed to create database directory: %w", err) + } + } + + var err error + UserDB, err = sql.Open("sqlite", dbPath) + if err != nil { + return fmt.Errorf("failed to open user database: %w", err) + } + + UserDB.SetMaxOpenConns(1) // SQLite serializes writes. + UserDB.SetMaxIdleConns(1) + UserDB.SetConnMaxLifetime(5 * time.Minute) + + if err := UserDB.Ping(); err != nil { + return fmt.Errorf("failed to ping user database: %w", err) + } + + if err := initUserTable(); err != nil { + return fmt.Errorf("failed to initialize user table: %w", err) + } + + return nil +} + +func initUserTable() error { + createTableSQL := ` + CREATE TABLE IF NOT EXISTS users ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + username TEXT NOT NULL UNIQUE, + password TEXT NOT NULL, + level TEXT NOT NULL DEFAULT 'admin', + register_date TEXT NOT NULL + );` + _, err := UserDB.Exec(createTableSQL) + return err +} + +// HashPassword returns the SHA-256 hex hash of a password. +func HashPassword(password string) string { + hash := sha256.Sum256([]byte(password)) + return hex.EncodeToString(hash[:]) +} + +// CreateUser inserts a new user into the database. +func CreateUser(username, password, level string) (int64, error) { + hashedPassword := HashPassword(password) + registerDate := time.Now().Format("2006-01-02 15:04:05") + + result, err := UserDB.Exec( + "INSERT INTO users (username, password, level, register_date) VALUES (?, ?, ?, ?)", + username, hashedPassword, level, registerDate, + ) + if err != nil { + return 0, fmt.Errorf("failed to create user: %w", err) + } + + return result.LastInsertId() +} + +// GetUserByUsername retrieves a user by username and validates the password. +func GetUserByUsername(username, password string) (int64, string, string, string, error) { + hashedPassword := HashPassword(password) + + var id int64 + var dbUsername string + var level string + var registerDate string + + err := UserDB.QueryRow( + "SELECT id, username, level, register_date FROM users WHERE username = ? AND password = ?", + username, hashedPassword, + ).Scan(&id, &dbUsername, &level, ®isterDate) + + if err == sql.ErrNoRows { + return 0, "", "", "", fmt.Errorf("invalid username or password") + } + if err != nil { + return 0, "", "", "", fmt.Errorf("database error: %w", err) + } + + return id, dbUsername, level, registerDate, nil +} + +// GetUserCount returns the total number of users in the database. +func GetUserCount() (int, error) { + var count int + err := UserDB.QueryRow("SELECT COUNT(*) FROM users").Scan(&count) + return count, err +} + +// GetUserByID retrieves a user by their ID. +func GetUserByID(id int64) (string, string, string, error) { + var username string + var level string + var registerDate string + + err := UserDB.QueryRow( + "SELECT username, level, register_date FROM users WHERE id = ?", + id, + ).Scan(&username, &level, ®isterDate) + + if err == sql.ErrNoRows { + return "", "", "", fmt.Errorf("user not found") + } + if err != nil { + return "", "", "", fmt.Errorf("database error: %w", err) + } + + return username, level, registerDate, nil +} + +// CloseUserDB closes the user database connection. +func CloseUserDB() { + if UserDB != nil { + UserDB.Close() + } +} diff --git a/db/log.db b/db/log.db new file mode 100644 index 0000000..e69de29 diff --git a/db/user.db b/db/user.db new file mode 100644 index 0000000..c341f40 Binary files /dev/null and b/db/user.db differ diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..d15e8b6 --- /dev/null +++ b/go.mod @@ -0,0 +1,22 @@ +module nukumizu-backend + +go 1.25.0 + +require ( + github.com/gorilla/websocket v1.5.3 + gopkg.in/mail.v2 v2.3.1 + modernc.org/sqlite v1.55.0 +) + +require ( + github.com/dustin/go-humanize v1.0.1 // indirect + github.com/google/uuid v1.6.0 // indirect + github.com/mattn/go-isatty v0.0.20 // indirect + github.com/ncruces/go-strftime v1.0.0 // indirect + github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect + golang.org/x/sys v0.46.0 // indirect + gopkg.in/alexcesaro/quotedprintable.v3 v3.0.0-20150716171945-2caba252f4dc // indirect + modernc.org/libc v1.74.1 // indirect + modernc.org/mathutil v1.7.1 // indirect + modernc.org/memory v1.11.0 // indirect +) diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..0b921ef --- /dev/null +++ b/go.sum @@ -0,0 +1,57 @@ +github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY= +github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto= +github.com/google/pprof v0.0.0-20250317173921-a4b03ec1a45e h1:ijClszYn+mADRFY17kjQEVQ1XRhq2/JR1M3sGqeJoxs= +github.com/google/pprof v0.0.0-20250317173921-a4b03ec1a45e/go.mod h1:boTsfXsheKC2y+lKOCMpSfarhxDeIzfZG1jqGcPl3cA= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg= +github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= +github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k= +github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM= +github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= +github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= +github.com/ncruces/go-strftime v1.0.0 h1:HMFp8mLCTPp341M/ZnA4qaf7ZlsbTc+miZjCLOFAw7w= +github.com/ncruces/go-strftime v1.0.0/go.mod h1:Fwc5htZGVVkseilnfgOVb9mKy6w1naJmn9CehxcKcls= +github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec h1:W09IVJc94icq4NjY3clb7Lk8O1qJ8BdBEF8z0ibU0rE= +github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo= +golang.org/x/mod v0.37.0 h1:vF1DjpVEshcIqoEaauuHebaLk1O1forxjxBaVn884JQ= +golang.org/x/mod v0.37.0/go.mod h1:m8S8VeM9r4dzDwjrKO0a1sZP3YjeMamRRlD+fmR2Q/0= +golang.org/x/sync v0.21.0 h1:HLII4xRRTtCRkxYp4HNFF0Js/Og6q2i++KXbg0gHCwM= +golang.org/x/sync v0.21.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= +golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.46.0 h1:noSf2Fq6F8DBgS+LysIkx7rIExoNHJsxOAtPp4rthXw= +golang.org/x/sys v0.46.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/tools v0.47.0 h1:7Kn5x/d1svx/PzryTsqeoZN4TZwqeH5pGWjefhLi/1Q= +golang.org/x/tools v0.47.0/go.mod h1:dFHnyTvFWY212G+h7ZY4Vsp/K3U4/7W9TyVaAul8uCA= +gopkg.in/alexcesaro/quotedprintable.v3 v3.0.0-20150716171945-2caba252f4dc h1:2gGKlE2+asNV9m7xrywl36YYNnBG5ZQ0r/BOOxqPpmk= +gopkg.in/alexcesaro/quotedprintable.v3 v3.0.0-20150716171945-2caba252f4dc/go.mod h1:m7x9LTH6d71AHyAX77c9yqWCCa3UKHcVEj9y7hAtKDk= +gopkg.in/mail.v2 v2.3.1 h1:WYFn/oANrAGP2C0dcV6/pbkPzv8yGzqTjPmTeO7qoXk= +gopkg.in/mail.v2 v2.3.1/go.mod h1:htwXN1Qh09vZJ1NVKxQqHPBaCBbzKhp5GzuJEA4VJWw= +modernc.org/cc/v4 v4.29.0 h1:CXgwL8cvxmyzBQZzbSl/6xFtMCryb6u8IOqDci39cgc= +modernc.org/cc/v4 v4.29.0/go.mod h1:OnovgIhbbMXMu1aISnJ0wvVD1KnW+cAUJkIrAWh+kVI= +modernc.org/ccgo/v4 v4.34.6 h1:sBgfIwyN0TQ9C5hwIeuqyeAKyMWnbvj2fvpF4L11uzU= +modernc.org/ccgo/v4 v4.34.6/go.mod h1:SZ8YcN9NG7XVsQYdm6jYBvi8PQP1qi+kqB6OhjqI3Fk= +modernc.org/fileutil v1.4.0 h1:j6ZzNTftVS054gi281TyLjHPp6CPHr2KCxEXjEbD6SM= +modernc.org/fileutil v1.4.0/go.mod h1:EqdKFDxiByqxLk8ozOxObDSfcVOv/54xDs/DUHdvCUU= +modernc.org/gc/v2 v2.6.5 h1:nyqdV8q46KvTpZlsw66kWqwXRHdjIlJOhG6kxiV/9xI= +modernc.org/gc/v2 v2.6.5/go.mod h1:YgIahr1ypgfe7chRuJi2gD7DBQiKSLMPgBQe9oIiito= +modernc.org/gc/v3 v3.1.4 h1:2g65LGVSmFQrXeITAw97x7hCRvZFcyE1uDP+7Vng7JI= +modernc.org/gc/v3 v3.1.4/go.mod h1:HFK/6AGESC7Ex+EZJhJ2Gni6cTaYpSMmU/cT9RmlfYY= +modernc.org/goabi0 v0.2.0 h1:HvEowk7LxcPd0eq6mVOAEMai46V+i7Jrj13t4AzuNks= +modernc.org/goabi0 v0.2.0/go.mod h1:CEFRnnJhKvWT1c1JTI3Avm+tgOWbkOu5oPA8eH8LnMI= +modernc.org/libc v1.74.1 h1:bdR4VTKFMC4966QSNZ05XLGI/VwzVa2kTUX51Dm0riQ= +modernc.org/libc v1.74.1/go.mod h1:uH4t5bOx3G3g9Xcmj10YKlTcVISlRDwv8VoQJG9n8Os= +modernc.org/mathutil v1.7.1 h1:GCZVGXdaN8gTqB1Mf/usp1Y/hSqgI2vAGGP4jZMCxOU= +modernc.org/mathutil v1.7.1/go.mod h1:4p5IwJITfppl0G4sUEDtCr4DthTaT47/N3aT6MhfgJg= +modernc.org/memory v1.11.0 h1:o4QC8aMQzmcwCK3t3Ux/ZHmwFPzE6hf2Y5LbkRs+hbI= +modernc.org/memory v1.11.0/go.mod h1:/JP4VbVC+K5sU2wZi9bHoq2MAkCnrt2r98UGeSK7Mjw= +modernc.org/opt v0.2.0 h1:tGyef5ApycA7FSEOMraay9SaTk5zmbx7Tu+cJs4QKZg= +modernc.org/opt v0.2.0/go.mod h1:03fq9lsNfvkYSfxrfUhZCWPk1lm4cq4N+Bh//bEtgns= +modernc.org/sortutil v1.2.1 h1:+xyoGf15mM3NMlPDnFqrteY07klSFxLElE2PVuWIJ7w= +modernc.org/sortutil v1.2.1/go.mod h1:7ZI3a3REbai7gzCLcotuw9AC4VZVpYMjDzETGsSMqJE= +modernc.org/sqlite v1.55.0 h1:hIFh0MCH0rGinQ/4KYb5/UbCkRkb+UP+OkLCVWa5MTM= +modernc.org/sqlite v1.55.0/go.mod h1:4ntCLuNmnH8+GNqjka1wNg7KJd5/Hi5FYp8K+XQ7GZw= +modernc.org/strutil v1.2.1 h1:UneZBkQA+DX2Rp35KcM69cSsNES9ly8mQWD71HKlOA0= +modernc.org/strutil v1.2.1/go.mod h1:EHkiggD70koQxjVdSBM3JKM7k6L0FbGE5eymy9i3B9A= +modernc.org/token v1.1.0 h1:Xl7Ap9dKaEs5kLoOQeQmPWevfnk/DM5qcLcYlA8ys6Y= +modernc.org/token v1.1.0/go.mod h1:UGzOrNV1mAFSEB63lOFHIpNRUVMvYTc6yu1SMY/XTDM= diff --git a/handler/bot.go b/handler/bot.go new file mode 100644 index 0000000..8aded73 --- /dev/null +++ b/handler/bot.go @@ -0,0 +1,96 @@ +package handler + +import ( + "encoding/json" + "net/http" + + "nukumizu-backend/internal/controller" + "nukumizu-backend/postLog" + "nukumizu-backend/utils" +) + +// OneBotMessage represents a OneBot 11 message event forwarded from napcat-bridge. +type OneBotMessage 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 string `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"` +} + +// BotMessageHandler handles POST /api/bot/msg/recv. +// Receives OneBot 11 events forwarded from napcat-bridge and +// routes them to the appropriate bot controller. +func BotMessageHandler(w http.ResponseWriter, r *http.Request) { + if !utils.Auth(w, r, "POST", "bot") { + return + } + + var event OneBotMessage + if err := json.NewDecoder(r.Body).Decode(&event); err != nil { + postLog.Debug("Failed to parse bot message: " + err.Error()) + utils.SendErrorResponse(w, http.StatusBadRequest, "invalid message format") + return + } + + // Only handle message events. + if event.PostType != "message" { + utils.SendSuccessResponse(w, "event ignored", map[string]interface{}{ + "post_type": event.PostType, + }) + return + } + + // Route to the Napcat/QQ controller. + ctrl := controller.GetController("qq(napcat)") + if ctrl == nil { + postLog.Debug("QQ controller not available, ignoring message") + utils.SendSuccessResponse(w, "controller not available", nil) + return + } + + cmdCtrl, ok := ctrl.(controller.CommandController) + if !ok { + utils.SendErrorResponse(w, http.StatusInternalServerError, "controller does not support commands") + return + } + + // Determine chat information. + chatID := event.GroupID + chatType := event.MessageType + if chatType == "private" { + chatID = event.UserID + } + + // Build a command from the raw message. + cmd := controller.Command{ + RawText: event.RawMessage, + ChatID: chatID, + ChatType: chatType, + SenderID: event.UserID, + } + + response, err := cmdCtrl.HandleCommand(cmd) + if err != nil { + postLog.Error("Bot command handling failed: " + err.Error()) + utils.SendErrorResponse(w, http.StatusInternalServerError, "command handling failed: "+err.Error()) + return + } + + if response != "" { + utils.SendSuccessResponse(w, "", map[string]interface{}{ + "response": response, + "chatID": chatID, + "chatType": chatType, + }) + } else { + utils.SendSuccessResponse(w, "no response", nil) + } +} diff --git a/handler/health.go b/handler/health.go new file mode 100644 index 0000000..90dc63a --- /dev/null +++ b/handler/health.go @@ -0,0 +1,28 @@ +package handler + +import ( + "net/http" + + "nukumizu-backend/database" + "nukumizu-backend/utils" +) + +// HealthHandler handles GET /health. +func HealthHandler(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodGet { + utils.SendErrorResponse(w, http.StatusMethodNotAllowed, "method not allowed") + return + } + + dbStatus := "ok" + if database.UserDB == nil { + dbStatus = "not initialized" + } else if err := database.UserDB.Ping(); err != nil { + dbStatus = "error: " + err.Error() + } + + utils.SendSuccessResponse(w, "", map[string]interface{}{ + "status": "ok", + "database": dbStatus, + }) +} diff --git a/handler/server.go b/handler/server.go new file mode 100644 index 0000000..9ca4ae8 --- /dev/null +++ b/handler/server.go @@ -0,0 +1,115 @@ +package handler + +import ( + "encoding/json" + "fmt" + "net/http" + + "nukumizu-backend/internal/komari" + "nukumizu-backend/internal/node" + "nukumizu-backend/internal/template" + "nukumizu-backend/postLog" + "nukumizu-backend/utils" +) + +// ServerExecRequest represents the request body for POST /api/server/exec. +type ServerExecRequest struct { + UUID []string `json:"uuid"` + Command string `json:"command"` +} + +// ServerListHandler handles GET /api/server/list. +func ServerListHandler(w http.ResponseWriter, r *http.Request) { + if !utils.Auth(w, r, "GET", "bot") { + return + } + + tracker := node.GetTracker() + if tracker == nil { + utils.SendErrorResponse(w, http.StatusInternalServerError, "node tracker not initialized") + return + } + + params := template.BuildParamsFromServerList() + result := template.Render("", params) + + utils.SendSuccessResponse(w, "", map[string]interface{}{ + "list": result, + }) +} + +// ServerGetStatusHandler handles GET /api/server/getStatus?uuid=xxx. +func ServerGetStatusHandler(w http.ResponseWriter, r *http.Request) { + if !utils.Auth(w, r, "GET", "bot") { + return + } + + uuid := r.URL.Query().Get("uuid") + if uuid == "" { + utils.SendErrorResponse(w, http.StatusBadRequest, "missing uuid parameter") + return + } + + // First try to get data from the local tracker. + tracker := node.GetTracker() + if tracker != nil { + if n, exists := tracker.GetNode(uuid); exists && n.LatestReport != nil { + utils.SendSuccessResponse(w, "", map[string]interface{}{ + "uuid": uuid, + "report": n.LatestReport, + }) + return + } + } + + utils.SendErrorResponse(w, http.StatusNotFound, fmt.Sprintf("no recent data for uuid: %s", uuid)) +} + +// ServerExecHandler handles POST /api/server/exec. +func ServerExecHandler(w http.ResponseWriter, r *http.Request) { + if !utils.Auth(w, r, "POST", "bot") { + return + } + + var req ServerExecRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + utils.SendErrorResponse(w, http.StatusBadRequest, "invalid request body") + return + } + + if len(req.UUID) == 0 { + utils.SendErrorResponse(w, http.StatusBadRequest, "uuid array is required") + return + } + + if req.Command == "" { + utils.SendErrorResponse(w, http.StatusBadRequest, "command is required") + return + } + + // Get the Komari client from the global state. + komariClient := komari.GetClient() + if komariClient == nil { + utils.SendErrorResponse(w, http.StatusInternalServerError, "komari client not initialized") + return + } + + taskID, err := komariClient.ExecTask(req.UUID, req.Command) + if err != nil { + postLog.Error("Komari exec task failed: " + err.Error()) + utils.SendErrorResponse(w, http.StatusInternalServerError, "failed to execute command: "+err.Error()) + return + } + + results, err := komariClient.PollTaskResult(taskID) + if err != nil { + postLog.Error("Komari task polling failed: " + err.Error()) + utils.SendErrorResponse(w, http.StatusInternalServerError, "failed to get task result: "+err.Error()) + return + } + + utils.SendSuccessResponse(w, "", map[string]interface{}{ + "taskID": taskID, + "results": results, + }) +} diff --git a/handler/user.go b/handler/user.go new file mode 100644 index 0000000..7792de3 --- /dev/null +++ b/handler/user.go @@ -0,0 +1,124 @@ +package handler + +import ( + "encoding/json" + "net/http" + + db "nukumizu-backend/database" + "nukumizu-backend/postLog" + "nukumizu-backend/utils" +) + +// UserLoginRequest represents the login request body. +type UserLoginRequest struct { + Username string `json:"username"` + Password string `json:"password"` +} + +// UserLoginHandler handles POST /api/user/login. +func UserLoginHandler(w http.ResponseWriter, r *http.Request) { + if !utils.Auth(w, r, "POST", "None") { + return + } + + var req UserLoginRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + utils.SendErrorResponse(w, http.StatusBadRequest, "invalid request body") + return + } + + if req.Username == "" || req.Password == "" { + utils.SendErrorResponse(w, http.StatusBadRequest, "username and password are required") + return + } + + userID, username, level, registerDate, err := db.GetUserByUsername(req.Username, req.Password) + if err != nil { + postLog.Warning("Login failed for user " + req.Username + ": " + err.Error()) + utils.SendErrorResponse(w, http.StatusUnauthorized, "invalid username or password") + return + } + + token, err := utils.GenerateToken() + if err != nil { + postLog.Error("Failed to generate token: " + err.Error()) + utils.SendErrorResponse(w, http.StatusInternalServerError, "failed to generate token") + return + } + + utils.AddToken(token, userID, level, username) + postLog.Info("User logged in: " + username) + + utils.SendSuccessResponse(w, "login successful", map[string]interface{}{ + "token": token, + "userID": userID, + "username": username, + "level": level, + "registerDate": registerDate, + }) +} + +// UserRegisterHandler handles POST /api/user/register. +// Only allowed when there are no existing users in the database. +func UserRegisterHandler(w http.ResponseWriter, r *http.Request) { + if !utils.Auth(w, r, "POST", "None") { + return + } + + // Check if any user already exists. + count, err := db.GetUserCount() + if err != nil { + postLog.Error("Failed to check user count: " + err.Error()) + utils.SendErrorResponse(w, http.StatusInternalServerError, "database error") + return + } + if count > 0 { + utils.SendErrorResponse(w, http.StatusForbidden, "registration is closed: users already exist") + return + } + + var req UserLoginRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + utils.SendErrorResponse(w, http.StatusBadRequest, "invalid request body") + return + } + + if req.Username == "" || req.Password == "" { + utils.SendErrorResponse(w, http.StatusBadRequest, "username and password are required") + return + } + + if len(req.Username) > 50 { + utils.SendErrorResponse(w, http.StatusBadRequest, "username too long (max 50 characters)") + return + } + + if len(req.Password) < 6 || len(req.Password) > 100 { + utils.SendErrorResponse(w, http.StatusBadRequest, "password must be between 6 and 100 characters") + return + } + + userID, err := db.CreateUser(req.Username, req.Password, "admin") + if err != nil { + postLog.Error("Failed to register user: " + err.Error()) + utils.SendErrorResponse(w, http.StatusInternalServerError, "failed to register user, username may already exist") + return + } + + token, err := utils.GenerateToken() + if err != nil { + postLog.Error("Failed to generate token: " + err.Error()) + utils.SendErrorResponse(w, http.StatusInternalServerError, "failed to generate token") + return + } + + utils.AddToken(token, userID, "admin", req.Username) + postLog.Info("User registered: " + req.Username) + + utils.SendSuccessResponse(w, "user registered successfully", map[string]interface{}{ + "token": token, + "userID": userID, + "username": req.Username, + "level": "admin", + }) +} diff --git a/internal/controller/controller.go b/internal/controller/controller.go new file mode 100644 index 0000000..50c1133 --- /dev/null +++ b/internal/controller/controller.go @@ -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 + } +} diff --git a/internal/controller/email.go b/internal/controller/email.go new file mode 100644 index 0000000..9c86eca --- /dev/null +++ b/internal/controller/email.go @@ -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 +} diff --git a/internal/controller/ntfy.go b/internal/controller/ntfy.go new file mode 100644 index 0000000..0cdb57f --- /dev/null +++ b/internal/controller/ntfy.go @@ -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 +} diff --git a/internal/controller/qq.go b/internal/controller/qq.go new file mode 100644 index 0000000..38a26c8 --- /dev/null +++ b/internal/controller/qq.go @@ -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 ", 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 ", 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 ", 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 ", 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() +} diff --git a/internal/controller/telegram.go b/internal/controller/telegram.go new file mode 100644 index 0000000..3356ad8 --- /dev/null +++ b/internal/controller/telegram.go @@ -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 ", 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 ", 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 ", 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 ", 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 +} diff --git a/internal/controller/webhook.go b/internal/controller/webhook.go new file mode 100644 index 0000000..f0c956c --- /dev/null +++ b/internal/controller/webhook.go @@ -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 +} diff --git a/internal/komari/client.go b/internal/komari/client.go new file mode 100644 index 0000000..fd3416d --- /dev/null +++ b/internal/komari/client.go @@ -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 +} diff --git a/internal/komari/ws.go b/internal/komari/ws.go new file mode 100644 index 0000000..aaf9f0c --- /dev/null +++ b/internal/komari/ws.go @@ -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 +} diff --git a/internal/node/tracker.go b/internal/node/tracker.go new file mode 100644 index 0000000..01b39ce --- /dev/null +++ b/internal/node/tracker.go @@ -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 +} + diff --git a/internal/template/template.go b/internal/template/template.go new file mode 100644 index 0000000..0d01d81 --- /dev/null +++ b/internal/template/template.go @@ -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") +} diff --git a/logs.db/log.db b/logs.db/log.db new file mode 100644 index 0000000..e69de29 diff --git a/logs.db/user.db b/logs.db/user.db new file mode 100644 index 0000000..c341f40 Binary files /dev/null and b/logs.db/user.db differ diff --git a/main.go b/main.go new file mode 100644 index 0000000..156f41f --- /dev/null +++ b/main.go @@ -0,0 +1,294 @@ +package main + +import ( + "flag" + "fmt" + "log" + "net/http" + "os" + "os/signal" + "syscall" + "time" + + "nukumizu-backend/config" + "nukumizu-backend/database" + "nukumizu-backend/internal/controller" + "nukumizu-backend/internal/komari" + "nukumizu-backend/internal/node" + "nukumizu-backend/postLog" + "nukumizu-backend/utils" +) + +// SoftwareInfo holds build metadata. +type SoftwareInfo struct { + Name string + Version string + Developer string + BuildVer int16 + Description string + BuildType string +} + +var softwareInfo = SoftwareInfo{ + Name: "Nukumizu", + Version: "0.1.0", + Developer: "Madobi Nanami", + BuildVer: 1, + Description: "Remote server monitoring and command execution subsystem for Komari", + BuildType: "Debug", +} + +func main() { + // Parse CLI flags. + configPath := flag.String("config", "config.json", "Path to configuration file") + flag.Parse() + + // Log startup banner. + postLog.Info(fmt.Sprintf("%s Ver.%s.%d.%s", softwareInfo.Name, softwareInfo.Version, softwareInfo.BuildVer, softwareInfo.BuildType)) + postLog.Info(fmt.Sprintf("Developed by %s", softwareInfo.Developer)) + postLog.Info(softwareInfo.Description) + + // Load configuration. + cfg, err := config.LoadConfig(*configPath) + if err != nil { + log.Fatalf("Failed to load config: %v", err) + } + + // Initialize logging. + postLog.SetDebugMode(cfg.System.DebugMode) + postLog.InitLogBroadcaster() + + dbPath := cfg.DBPath + if dbPath == "" { + dbPath = "./db" + } + + if err := postLog.InitLogsDatabase(fmt.Sprintf("%s/log.db", dbPath)); err != nil { + postLog.Fatal("Failed to initialize logs database: " + err.Error()) + } + + // Initialize database. + if err := database.InitUserDB(fmt.Sprintf("%s/user.db", dbPath)); err != nil { + postLog.Fatal("Failed to initialize user database: " + err.Error()) + } + defer database.CloseUserDB() + + // Initialize node tracker. + node.InitTracker() + + // Initialize controller manager. + controller.InitManager() + + // Initialize rate limiter. + utils.InitRateLimiter(100, time.Minute) + + // Start token cleaner. + utils.StartTokenCleaner() + + postLog.Info("Starting Nukumizu server...") + + // --- Login to Komari --- + komari.InitClient(cfg.Komari.DashboardURL) + + if err := komari.LoginAndStart(); err != nil { + postLog.Fatal("Failed to login to Komari: " + err.Error()) + } + + // --- Start Komari WebSocket connection --- + komari.InitWSClient(cfg.Komari.DashboardURL, 5) + wsClient := komari.GetWSClient() + if wsClient != nil { + wsClient.SetOnReconnectFail(func() { + postLog.Error("Komari WebSocket reconnection exhausted") + controller.GetManager().NotifyAllAdmins("Komari WebSocket connection lost after 5 retry attempts") + }) + + if err := wsClient.Connect(); err != nil { + postLog.Error("Failed to connect Komari WebSocket: " + err.Error()) + } + } + + // --- Initialize controllers --- + initControllers() + + // --- Start background tasks --- + startBackgroundTasks(cfg) + + // --- Setup HTTP router --- + mux := SetupRouter() + + // Apply middleware chain. + handler := utils.RateLimitMiddleware(mux) + handler = utils.CORSMiddleware(handler) + handler = utils.XSSProtectionMiddleware(handler) + + addr := fmt.Sprintf("%s:%s", cfg.System.ListenAddr, cfg.System.ListenPort) + postLog.Info(fmt.Sprintf("Server listening on %s", addr)) + + // Start HTTP server in a goroutine. + go func() { + if err := http.ListenAndServe(addr, handler); err != nil { + postLog.Fatal("Server error: " + err.Error()) + } + }() + + // --- Graceful shutdown --- + quit := make(chan os.Signal, 1) + signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM) + <-quit + + postLog.Info("Shutting down Nukumizu server...") + + // Stop WebSocket client. + if wsClient != nil { + wsClient.Stop() + } + + // Stop all controllers. + mgr := controller.GetManager() + if mgr != nil { + mgr.StopAll() + } + + postLog.Info("Server stopped") +} + +// initControllers initializes and starts all configured controllers. +func initControllers() { + cfg := config.GetConfig() + mgr := controller.GetManager() + if mgr == nil { + return + } + + // QQ (Napcat) controller. + qqCtrl := controller.NewQQController(cfg.ControllerMethod.QQ) + mgr.Register(qqCtrl) + go func() { + defer func() { + if r := recover(); r != nil { + postLog.Error(fmt.Sprintf("QQ controller panic: %v", r)) + } + }() + if err := qqCtrl.Start(); err != nil { + postLog.Error("Failed to start QQ controller: " + err.Error()) + } + }() + + // Telegram controller. + tgCtrl := controller.NewTelegramController(cfg.ControllerMethod.Telegram) + mgr.Register(tgCtrl) + go func() { + defer func() { + if r := recover(); r != nil { + postLog.Error(fmt.Sprintf("Telegram controller panic: %v", r)) + } + }() + if err := tgCtrl.Start(); err != nil { + postLog.Error("Failed to start Telegram controller: " + err.Error()) + } + }() + + // Email controller (status-only). + emailCtrl := controller.NewEmailController(cfg.ControllerMethod.Email) + mgr.Register(emailCtrl) + go func() { + defer func() { + if r := recover(); r != nil { + postLog.Error(fmt.Sprintf("Email controller panic: %v", r)) + } + }() + if err := emailCtrl.Start(); err != nil { + postLog.Error("Failed to start Email controller: " + err.Error()) + } + }() + + // Ntfy controller (status-only). + ntfyCtrl := controller.NewNtfyController(cfg.ControllerMethod.Ntfy) + mgr.Register(ntfyCtrl) + go func() { + defer func() { + if r := recover(); r != nil { + postLog.Error(fmt.Sprintf("Ntfy controller panic: %v", r)) + } + }() + if err := ntfyCtrl.Start(); err != nil { + postLog.Error("Failed to start Ntfy controller: " + err.Error()) + } + }() + + // Webhook controller (status-only). + webhookCtrl := controller.NewWebhookController(cfg.ControllerMethod.Webhook) + mgr.Register(webhookCtrl) + go func() { + defer func() { + if r := recover(); r != nil { + postLog.Error(fmt.Sprintf("Webhook controller panic: %v", r)) + } + }() + if err := webhookCtrl.Start(); err != nil { + postLog.Error("Failed to start Webhook controller: " + err.Error()) + } + }() +} + +// startBackgroundTasks starts periodic background goroutines. +func startBackgroundTasks(_ *config.Config) { + // Refresh node list every 5 minutes. + go func() { + defer func() { + if r := recover(); r != nil { + postLog.Error(fmt.Sprintf("Node refresh panic: %v", r)) + } + }() + ticker := time.NewTicker(5 * time.Minute) + defer ticker.Stop() + for range ticker.C { + client := komari.GetClient() + if client == nil { + continue + } + nodes, err := client.FetchNodes() + if err != nil { + postLog.Warning("Failed to refresh node list: " + err.Error()) + continue + } + 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) + } + } + }() + + // Re-login to Komari every 12 hours. + go func() { + defer func() { + if r := recover(); r != nil { + postLog.Error(fmt.Sprintf("Komari re-login panic: %v", r)) + } + }() + ticker := time.NewTicker(12 * time.Hour) + defer ticker.Stop() + for range ticker.C { + postLog.Info("Performing scheduled Komari re-login...") + if err := komari.LoginAndStart(); err != nil { + postLog.Error("Komari re-login failed: " + err.Error()) + } + } + }() + + // Wire node status change notifications to all controllers. + tracker := node.GetTracker() + if tracker != nil { + tracker.OnStatusChange(func(change node.StatusChange) { + mgr := controller.GetManager() + if mgr != nil { + mgr.NotifyStatusChange(change) + } + }) + } +} diff --git a/router.go b/router.go new file mode 100644 index 0000000..823649d --- /dev/null +++ b/router.go @@ -0,0 +1,52 @@ +package main + +import ( + "fmt" + "net/http" + + "nukumizu-backend/handler" + "nukumizu-backend/postLog" +) + +// SetupRouter registers all HTTP routes and returns a configured ServeMux. +func SetupRouter() *http.ServeMux { + postLog.Info("Setting up routers...") + + mux := http.NewServeMux() + + // User endpoints. + mux.HandleFunc("/api/user/login", handler.UserLoginHandler) + mux.HandleFunc("/api/user/register", handler.UserRegisterHandler) + + // Server endpoints (authenticated). + mux.HandleFunc("/api/server/list", handler.ServerListHandler) + mux.HandleFunc("/api/server/getStatus", handler.ServerGetStatusHandler) + mux.HandleFunc("/api/server/exec", handler.ServerExecHandler) + + // Bot message receive endpoint (from napcat-bridge). + mux.HandleFunc("/api/bot/msg/recv", handler.BotMessageHandler) + + // Health check endpoint. + mux.HandleFunc("/health", handler.HealthHandler) + + // WebSocket log streaming endpoint. + logBroadcaster := postLog.GetLogBroadcaster() + if logBroadcaster != nil { + logSocketHandler := postLog.NewLogSocketHandler(logBroadcaster) + mux.HandleFunc("/api/system/getLogs", logSocketHandler.Handle) + } + + // Catch-all 404 handler. + mux.HandleFunc("/", NotFoundHandler) + + postLog.Info("Router setup completed") + return mux +} + +// NotFoundHandler returns a 404 JSON response for unknown routes. +func NotFoundHandler(w http.ResponseWriter, r *http.Request) { + postLog.Debug(fmt.Sprintf("Unknown request: %s %s", r.Method, r.URL.Path)) + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusNotFound) + fmt.Fprint(w, `{"success":false,"message":"route not found"}`) +} diff --git a/utils/auth.go b/utils/auth.go new file mode 100644 index 0000000..57f4add --- /dev/null +++ b/utils/auth.go @@ -0,0 +1,229 @@ +package utils + +import ( + "crypto/rand" + "encoding/hex" + "encoding/json" + "net/http" + "strconv" + "sync" + "time" + + "nukumizu-backend/config" + "nukumizu-backend/postLog" +) + +// TokenInfo holds information about an authenticated session token. +type TokenInfo struct { + UserID int64 `json:"userID"` + Level string `json:"level"` + UserName string `json:"userName"` + CreatedAt time.Time `json:"createdAt"` + LastAccess time.Time `json:"lastAccess"` +} + +var ( + tokenStore = make(map[string]*TokenInfo) + tokenStoreLock sync.RWMutex +) + +// GenerateToken creates a cryptographically random 32-byte hex token. +func GenerateToken() (string, error) { + bytes := make([]byte, 32) + if _, err := rand.Read(bytes); err != nil { + return "", err + } + return hex.EncodeToString(bytes), nil +} + +// AddToken adds a token to the in-memory token store. +func AddToken(token string, userID int64, level string, userName string) { + tokenStoreLock.Lock() + defer tokenStoreLock.Unlock() + now := time.Now() + tokenStore[token] = &TokenInfo{ + UserID: userID, + Level: level, + UserName: userName, + CreatedAt: now, + LastAccess: now, + } +} + +// GetTokenInfo retrieves token information from the store. +func GetTokenInfo(token string) (*TokenInfo, bool) { + tokenStoreLock.RLock() + defer tokenStoreLock.RUnlock() + info, exists := tokenStore[token] + return info, exists +} + +// RefreshToken updates the LastAccess time for an active token. +func RefreshToken(token string) { + tokenStoreLock.Lock() + defer tokenStoreLock.Unlock() + if info, exists := tokenStore[token]; exists { + info.LastAccess = time.Now() + } +} + +// RemoveToken deletes a token from the store. +func RemoveToken(token string) { + tokenStoreLock.Lock() + defer tokenStoreLock.Unlock() + delete(tokenStore, token) +} + +// GetUserIDFromRequest extracts the user ID from the request's X-Token header. +func GetUserIDFromRequest(r *http.Request) int64 { + token := r.Header.Get("X-Token") + if token == "" { + return 0 + } + tokenInfo, exists := GetTokenInfo(token) + if !exists { + return 0 + } + return tokenInfo.UserID +} + +// GetUserLevelFromRequest extracts the user level from the request's X-Token header. +func GetUserLevelFromRequest(r *http.Request) string { + token := r.Header.Get("X-Token") + if token == "" { + return "" + } + tokenInfo, exists := GetTokenInfo(token) + if !exists { + return "" + } + return tokenInfo.Level +} + +// Auth is the central authentication and authorization function. +// It validates the request method, X-Timestamp header (30min tolerance), +// X-Token header, and permission level. Returns true if the request is authorized. +// +// Permission levels: "None" (public, no token required), "bot", "admin". +// When level is "bot", both "bot" and "admin" tokens are accepted. +// When level is "admin", only "admin" tokens are accepted. +func Auth(w http.ResponseWriter, r *http.Request, targetMethod string, targetLevel string) bool { + // Validate HTTP method. + if r.Method != targetMethod { + SendErrorResponse(w, http.StatusMethodNotAllowed, "method not allowed") + return false + } + + // Validate X-Timestamp. + timestamp := r.Header.Get("X-Timestamp") + if !config.IsDebugMode() { + if timestamp == "" { + SendErrorResponse(w, http.StatusUnauthorized, "missing timestamp") + return false + } + + ts, err := strconv.ParseInt(timestamp, 10, 64) + if err != nil { + SendErrorResponse(w, http.StatusUnauthorized, "invalid timestamp") + return false + } + + now := time.Now().Unix() + diff := now - ts + if diff < 0 { + diff = -diff + } + // 30 minute tolerance per agent.md. + if diff > 1800 { + SendErrorResponse(w, http.StatusUnauthorized, "request expired") + return false + } + } + + // Public endpoints require no token. + if targetLevel == "None" { + return true + } + + // Validate X-Token. + token := r.Header.Get("X-Token") + if token == "" { + SendErrorResponse(w, http.StatusUnauthorized, "missing token") + return false + } + + tokenInfo, exists := GetTokenInfo(token) + if !exists { + SendErrorResponse(w, http.StatusUnauthorized, "invalid token") + return false + } + + // Check permission level. + // "bot" level accepts both "bot" and "admin" tokens. + // "admin" level accepts only "admin" tokens. + switch targetLevel { + case "admin": + if tokenInfo.Level != "admin" { + SendErrorResponse(w, http.StatusForbidden, "permission denied") + return false + } + case "bot": + if tokenInfo.Level != "bot" && tokenInfo.Level != "admin" { + SendErrorResponse(w, http.StatusForbidden, "permission denied") + return false + } + } + + // Refresh token last access time. + RefreshToken(token) + return true +} + +// CleanExpiredTokens removes tokens that have been idle for over 1 hour. +func CleanExpiredTokens() { + tokenStoreLock.Lock() + defer tokenStoreLock.Unlock() + now := time.Now() + for token, info := range tokenStore { + if now.Sub(info.LastAccess) > 1*time.Hour { + delete(tokenStore, token) + postLog.Debug("Expired token removed for user: " + info.UserName) + } + } +} + +// StartTokenCleaner starts a background goroutine that periodically cleans +// expired tokens. +func StartTokenCleaner() { + ticker := time.NewTicker(10 * time.Minute) + go func() { + for range ticker.C { + CleanExpiredTokens() + } + }() +} + +// SendSuccessResponse sends a standardized JSON success response. +func SendSuccessResponse(w http.ResponseWriter, message string, data map[string]interface{}) { + w.Header().Set("Content-Type", "application/json") + resp := map[string]interface{}{ + "success": true, + } + if message != "" { + resp["message"] = message + } + for k, v := range data { + resp[k] = v + } + json.NewEncoder(w).Encode(resp) +} + +// SendErrorResponse sends a standardized JSON error response. +func SendErrorResponse(w http.ResponseWriter, statusCode int, message string) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(statusCode) + json.NewEncoder(w).Encode(map[string]interface{}{ + "success": false, + "message": message, + }) +} diff --git a/utils/middleware.go b/utils/middleware.go new file mode 100644 index 0000000..8866130 --- /dev/null +++ b/utils/middleware.go @@ -0,0 +1,114 @@ +package utils + +import ( + "net/http" + "sync" + "time" +) + +// RateLimiter implements a simple per-IP rate limiter using a token bucket approach. +type RateLimiter struct { + mu sync.Mutex + visitors map[string]*visitor + rate int + interval time.Duration +} + +type visitor struct { + count int + lastSeen time.Time +} + +var globalLimiter *RateLimiter + +// InitRateLimiter initializes the global rate limiter. +// rate is the maximum number of requests per interval. +func InitRateLimiter(rate int, interval time.Duration) { + globalLimiter = &RateLimiter{ + visitors: make(map[string]*visitor), + rate: rate, + interval: interval, + } + // Start a background cleaner. + go func() { + for { + time.Sleep(interval) + globalLimiter.cleanup() + } + }() +} + +func (rl *RateLimiter) cleanup() { + rl.mu.Lock() + defer rl.mu.Unlock() + for ip, v := range rl.visitors { + if time.Since(v.lastSeen) > rl.interval { + delete(rl.visitors, ip) + } + } +} + +func (rl *RateLimiter) allow(ip string) bool { + rl.mu.Lock() + defer rl.mu.Unlock() + v, exists := rl.visitors[ip] + if !exists { + rl.visitors[ip] = &visitor{count: 1, lastSeen: time.Now()} + return true + } + if time.Since(v.lastSeen) > rl.interval { + v.count = 1 + v.lastSeen = time.Now() + return true + } + if v.count >= rl.rate { + return false + } + v.count++ + return true +} + +// RateLimitMiddleware limits the number of requests per IP address. +func RateLimitMiddleware(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if globalLimiter == nil { + next.ServeHTTP(w, r) + return + } + ip := r.RemoteAddr + if !globalLimiter.allow(ip) { + SendErrorResponse(w, http.StatusTooManyRequests, "rate limit exceeded") + return + } + next.ServeHTTP(w, r) + }) +} + +// CORSMiddleware sets CORS headers and handles preflight OPTIONS requests. +func CORSMiddleware(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Access-Control-Allow-Origin", "*") + w.Header().Set("Access-Control-Allow-Methods", "GET, POST, PUT, DELETE, OPTIONS") + w.Header().Set("Access-Control-Allow-Headers", "Content-Type, X-Token, X-Timestamp, Authorization") + w.Header().Set("Access-Control-Max-Age", "86400") + + if r.Method == http.MethodOptions { + w.WriteHeader(http.StatusNoContent) + return + } + + next.ServeHTTP(w, r) + }) +} + +// XSSProtectionMiddleware sets security headers to help prevent XSS attacks. +func XSSProtectionMiddleware(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("X-XSS-Protection", "1; mode=block") + w.Header().Set("X-Content-Type-Options", "nosniff") + w.Header().Set("X-Frame-Options", "DENY") + w.Header().Set("Referrer-Policy", "strict-origin-when-cross-origin") + w.Header().Set("Content-Security-Policy", "default-src 'self'; script-src 'self'; style-src 'self' 'unsafe-inline'") + next.ServeHTTP(w, r) + }) +}