Files
Nukumizu/main.go
T
NanamiAdmin aa7e7f87e1
Build / ubuntu-latest (push) Failing after 3m17s
Build / windows-latest (push) Canceled after 8m6s
feat(node): add GetNodeCount to get node count and remove no need node update log, only log when node count changed
2026-09-29 21:43:46 +08:00

322 lines
9.6 KiB
Go

package main
import (
"fmt"
"log"
"net/http"
"os"
"os/signal"
"syscall"
"time"
"nukumizu-backend/config"
"nukumizu-backend/database"
"nukumizu-backend/global"
"nukumizu-backend/internal/controller"
"nukumizu-backend/internal/controller/pipes"
"nukumizu-backend/internal/controller/pipes/qq_napcat"
"nukumizu-backend/internal/controller/pipes/telegram"
"nukumizu-backend/internal/komari"
"nukumizu-backend/internal/node"
"nukumizu-backend/postLog"
"nukumizu-backend/utils"
)
// CommitHash and BuildTime are set at build time using -ldflags.
var (
CommitHash string
BuildTime string
)
func main() {
// Set global software info.
global.SoftwareInfo.CommitHash = CommitHash
global.SoftwareInfo.BuildTime = BuildTime
// Log startup banner.
postLog.Info(fmt.Sprintf("%s Ver.%s.%d.%s.%s Developed by %s at %s", global.SoftwareInfo.Name, global.SoftwareInfo.Version, global.SoftwareInfo.BuildVer, global.SoftwareInfo.BuildType, global.SoftwareInfo.CommitHash, global.SoftwareInfo.Developer, global.SoftwareInfo.BuildTime))
// Load configuration.
cfg, err := config.LoadGlobalConfig(global.ConfigPath.Global)
if err != nil {
log.Fatalf("Failed to load global config: %v", err)
}
_, err = config.LoadBotUserConfig(global.ConfigPath.BotUserConfig)
if err != nil {
log.Fatalf("Failed to load bot user config: %v", err)
}
// Load the node registry config. Unlike the other files it is optional:
// missing or empty bot_node_config.json simply means every node keeps its
// default enableStatusNotify (true).
if err := config.LoadBotNodeConfig(global.ConfigPath.BotNodeConfig); err != nil {
log.Fatalf("Failed to load bot node config: %v", err)
}
// Initialize logging.
postLog.SetDebugMode(cfg.System.DebugMode)
postLog.InitLogBroadcaster()
// A settings update replaces the configuration in memory; these hooks push
// the new values into the state that was derived from the old one. The
// logger's debug flag is process-wide rather than read at every log call,
// and each controller holds its own copy of its channel settings plus the
// clients built from them.
config.OnReload(func(updated *config.Config) {
postLog.SetDebugMode(updated.System.DebugMode)
mgr := controller.GetManager()
if mgr == nil || !mgr.NeedsRebuild(updated.ControllerMethod) {
return
}
// Controller settings changed. Each controller reads its settings into
// fields when it is built, and the NapCat and Telegram ones own
// connections that cannot be re-pointed, so the change is applied by
// replacing the whole set rather than reconfiguring it in place.
postLog.Info("Controller settings changed; rebuilding every channel")
mgr.ReplaceAll(buildControllers(updated), updated.ControllerMethod)
})
dbPath := cfg.DBPath
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())
return
}
// --- 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()
// Send the bot initialization message to all enabled bot controllers
// (QQ_NapCat, Telegram), then the server list right after it once the
// initial status refresh has completed so the list reflects the real
// online/offline state.
if mgr := controller.GetManager(); mgr != nil {
mgr.ShowBotInitMessage()
if tracker := node.GetTracker(); tracker != nil {
tracker.OnBootstrapDone(func() {
if m := controller.GetManager(); m != nil {
m.ShowBotServerList()
}
})
}
}
// --- 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())
}
}()
// --- Start the incoming webhook listener ---
startWebhookServer(cfg)
// --- 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")
}
// startWebhookServer serves the incoming webhook API on its own listener. The
// API is not exposed on the main listener: external applications post alerts to
// this port only, so its rate limiter and CORS policy are configured
// independently. A failure to bind it is logged rather than fatal — the rest of
// the program (bots, status monitoring) keeps running without it.
func startWebhookServer(cfg *config.Config) {
if !cfg.Webhook.Enabled {
postLog.Warning("Incoming webhook API is disabled")
return
}
handler := utils.RateLimitMiddleware(SetupWebhookRouter())
handler = utils.CORSMiddleware(handler)
addr := fmt.Sprintf("%s:%s", cfg.Webhook.ListenAddr, cfg.Webhook.ListenPort)
postLog.Info(fmt.Sprintf("Webhook API listening on %s", addr))
go func() {
defer func() {
if r := recover(); r != nil {
postLog.Error(fmt.Sprintf("Webhook server panic: %v", r))
}
}()
if err := http.ListenAndServe(addr, handler); err != nil {
postLog.Error("Webhook server error: " + err.Error())
}
}()
}
// buildControllers constructs one controller per channel from the given
// configuration. The set is built fresh whenever controller settings change,
// because a controller reads its settings once at construction and two of them
// own connections that cannot be re-pointed.
func buildControllers(cfg *config.Config) []controller.Controller {
return []controller.Controller{
qq_napcat.NewQQController(cfg.ControllerMethod.QQ),
telegram.NewTelegramController(cfg.ControllerMethod.Telegram),
pipes.NewEmailController(cfg.ControllerMethod.Email),
pipes.NewNtfyController(cfg.ControllerMethod.Ntfy),
pipes.NewWebhookController(cfg.ControllerMethod.Webhook),
}
}
// initControllers installs the initial controller set.
func initControllers() {
cfg := config.Current()
mgr := controller.GetManager()
if mgr == nil || cfg == nil {
return
}
mgr.ReplaceAll(buildControllers(cfg), cfg.ControllerMethod)
}
// 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
}
orgNodeCount := node.GetTracker().GetNodeCount()
nodes, err := client.FetchNodes()
if err != nil {
postLog.Warning("Failed to refresh node list: " + err.Error())
continue
}
if node.GetTracker().GetNodeCount() != orgNodeCount {
postLog.Info(fmt.Sprintf("Refreshed node list: %d nodes", len(nodes)))
}
tracker := node.GetTracker()
if tracker != nil {
tracker.UpdateNodeList(komari.BuildNodeListData(nodes))
}
}
}()
// 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)
}
})
}
// Safety net for the initial refresh gate: if some node never reports a
// status, force-complete the initial refresh after a grace period so status
// notifications cannot stay blocked forever.
go func() {
defer func() {
if r := recover(); r != nil {
postLog.Error(fmt.Sprintf("Bootstrap timeout panic: %v", r))
}
}()
time.Sleep(2 * time.Minute) // Add 2 minutes grace period to allow nodes to report status.
if tracker := node.GetTracker(); tracker != nil {
tracker.CompleteBootstrap()
}
}()
}