package main import ( "context" "errors" "fmt" "log/slog" "net" "net/http" "os" "os/signal" "path/filepath" "strings" "syscall" "time" "git.asio.asia/nixevol/NixMsg/internal/config" "git.asio.asia/nixevol/NixMsg/internal/store" "git.asio.asia/nixevol/NixMsg/web" ) func cmdServe(_ []string) error { cfgPath := config.PathFromEnv() cfg, err := config.Load(cfgPath) if err != nil { return err } // P-WIRE-BEGIN if err := cfg.Validate(); err != nil { return err } setupJSONLogger(cfg.Log) // P-WIRE-END ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() return runServe(ctx, cfg) } func runServe(ctx context.Context, cfg config.Config) error { deps := wire() if err := os.MkdirAll(cfg.DataDir, 0o755); err != nil { return fmt.Errorf("mkdir data_dir: %w", err) } db, err := store.Open(cfg.DataDir, cfg.SQLiteSynchronous) if err != nil { return err } defer func() { _ = db.Close() }() // P-WIRE-BEGIN ok, err := store.HasAdminPassword(ctx, db.Write) if err != nil { return err } if !ok { return errors.New("admin password not initialized; run: nixmsg admin init") } // P-WIRE-END // 启动恢复入口已挂上(假实现为空操作);M 线替换 message.Service 后生效。 if recoverErr := deps.Messages.RecoverOnStart(ctx); recoverErr != nil { return fmt.Errorf("message recover: %w", recoverErr) } // 其余 deps 供后续 admin / broker / httpx 接线;此处显式引用避免未使用告警。 _ = deps.Identity _ = deps.Groups _ = deps.Presence _ = deps.Downlink _ = deps.Conns _ = deps.Uplink _ = deps.HashPool _ = deps.SessionTokens _ = deps.APITokens _ = deps.LoginLocks mux := http.NewServeMux() mux.HandleFunc("GET /healthz", func(w http.ResponseWriter, _ *http.Request) { w.WriteHeader(http.StatusOK) _, _ = w.Write([]byte("ok")) }) // P-WIRE-BEGIN mux.HandleFunc("GET /readyz", func(w http.ResponseWriter, r *http.Request) { if readyErr := db.Ready(r.Context()); readyErr != nil { slog.Error("readyz failed", "err", readyErr) http.Error(w, "not ready", http.StatusServiceUnavailable) return } w.WriteHeader(http.StatusOK) _, _ = w.Write([]byte("ok")) }) // P-WIRE-END // 确保前端资源被链接进二进制;完整静态托管由后续任务完善。 _ = web.Dist() ln, err := net.Listen("tcp", cfg.Listen) if err != nil { return fmt.Errorf("listen %s: %w", cfg.Listen, err) } if err := writeListenAddr(cfg.DataDir, ln.Addr().String()); err != nil { _ = ln.Close() return err } srv := &http.Server{ Handler: mux, ReadHeaderTimeout: 10 * time.Second, } errCh := make(chan error, 1) go func() { errCh <- srv.Serve(ln) }() select { case <-ctx.Done(): // P-WIRE-BEGIN // 先停止接受新连接,再等写队列最多 10 秒,然后断开并退出。 shutdownCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second) defer cancel() _ = srv.Shutdown(shutdownCtx) drainCtx, drainCancel := context.WithTimeout(context.Background(), 10*time.Second) defer drainCancel() if drainErr := db.Queue.Drain(drainCtx); drainErr != nil && !errors.Is(drainErr, context.DeadlineExceeded) { slog.Error("write queue drain", "err", drainErr) } // P-WIRE-END serveErr := <-errCh if serveErr != nil && !errors.Is(serveErr, http.ErrServerClosed) { return serveErr } return nil case serveErr := <-errCh: if errors.Is(serveErr, http.ErrServerClosed) { return nil } return serveErr } } func writeListenAddr(dataDir, addr string) error { path := filepath.Join(dataDir, "listen.addr") return os.WriteFile(path, []byte(addr+"\n"), 0o644) } // P-WIRE-BEGIN func setupJSONLogger(cfg config.LogConfig) { level := slog.LevelInfo switch strings.ToLower(strings.TrimSpace(cfg.Level)) { case "debug": level = slog.LevelDebug case "warn", "warning": level = slog.LevelWarn case "error": level = slog.LevelError } h := slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{Level: level}) slog.SetDefault(slog.New(h)) } // P-WIRE-END