// API 进程入口,负责装配依赖、启动 HTTP 服务和异步任务处理器并执行平滑关闭。 package main import ( "context" "errors" "log/slog" "net/http" "os" "os/signal" "syscall" "time" "juhe-factory/api/internal/cache" "juhe-factory/api/internal/config" "juhe-factory/api/internal/database" "juhe-factory/api/internal/datetime" "juhe-factory/api/internal/handler" adminmodule "juhe-factory/api/internal/modules/admin" "juhe-factory/api/internal/modules/productimage" "juhe-factory/api/internal/modules/prompt" "juhe-factory/api/internal/queue" "juhe-factory/api/internal/security" "juhe-factory/api/internal/server" "juhe-factory/api/internal/service" "juhe-factory/api/internal/storage" "juhe-factory/api/internal/worker" "github.com/hibiken/asynq" ) func main() { time.Local = datetime.Location cfg, err := config.Load() if err != nil { slog.Error("加载配置失败", "error", err) os.Exit(1) } gormDB, sqlDB, err := database.Open(cfg) if err != nil { slog.Error("连接 PostgreSQL 失败", "error", err) os.Exit(1) } defer sqlDB.Close() redisClient, err := cache.Open(cfg) if err != nil { slog.Error("连接 Redis 失败", "error", err) os.Exit(1) } defer redisClient.Close() queueClient := queue.NewClient(cfg) defer queueClient.Close() healthHandler := handler.NewHealth(sqlDB, redisClient) passwords := security.PasswordHasher{Time: cfg.Argon2Time, Memory: cfg.Argon2Memory, Parallelism: cfg.Argon2Parallelism, HashLength: cfg.Argon2HashLength, SaltLength: cfg.Argon2SaltLength} tokens := security.TokenService{Secret: []byte(cfg.JWTSecretKey), AccessTTL: time.Duration(cfg.JWTAccessMinutes) * time.Minute, RefreshTTL: time.Duration(cfg.JWTRefreshHours) * time.Hour} encryptor, err := security.NewEncryptor(cfg.EncryptionKey) if err != nil { slog.Error("加载敏感配置加密密钥失败", "error", err) os.Exit(1) } authService := service.NewAuth(gormDB, passwords, tokens) if err := authService.Bootstrap(cfg.BootstrapAdminUsername, cfg.BootstrapAdminPassword); err != nil { slog.Error("初始化管理员失败", "error", err) os.Exit(1) } cos, err := storage.NewCOS(context.Background(), cfg) if err != nil { slog.Error("初始化 COS 失败", "error", err) os.Exit(1) } adminDataService := &service.AdminData{DB: gormDB, Passwords: passwords, Encryptor: encryptor, HTTPClient: http.DefaultClient, EncryptionKeyVersion: cfg.EncryptionKeyVersion} webService := &service.Web{DB: gormDB, Passwords: passwords, Tokens: security.WebTokenService{Secret: []byte(cfg.JWTSecretKey), AccessTTL: time.Duration(cfg.JWTAccessMinutes) * time.Minute}, RefreshTTL: time.Duration(cfg.JWTRefreshHours) * time.Hour} creativeService := &service.Creative{DB: gormDB, Queue: queueClient, Encryptor: encryptor, HTTPClient: http.DefaultClient} adminService := adminmodule.NewService(gormDB) promptService := prompt.NewService(gormDB) productImageService := productimage.NewService(gormDB, queueClient) creativeHandler := handler.NewCreative(creativeService, promptService, cos) httpServer := server.New(cfg, healthHandler, handler.NewAdminAuth(authService), handler.NewAdminData(adminDataService, adminService, promptService, cos), handler.NewWeb(webService, cos), creativeHandler, handler.NewProductImage(productImageService, cos)) aiHTTP := worker.NewHTTPClient(cfg.AIHTTPMaxConnections) generationWorker := worker.NewGeneration(gormDB, queueClient, cos, encryptor, aiHTTP, time.Duration(cfg.AIPollIntervalSeconds)*time.Second) mediaWorker := worker.NewMedia(gormDB, cos, encryptor, aiHTTP, cfg.FFmpegPath, cfg.FFprobePath, cfg.PythonPath, cfg.ASRScriptPath) dramaParseWorker := worker.NewDramaParse(gormDB, queueClient, encryptor, aiHTTP) scriptAnalysisWorker := worker.NewScriptAnalysis(gormDB, queueClient, encryptor, aiHTTP) workerMux := asynq.NewServeMux() generationWorker.Register(workerMux) mediaWorker.Register(workerMux) dramaParseWorker.Register(workerMux) scriptAnalysisWorker.Register(workerMux) queueServer := queue.NewServer(cfg) workerErr := make(chan error, 1) go func() { slog.Info("AI 任务 Worker 已启动", "concurrency", cfg.AIWorkerConcurrency) workerErr <- queueServer.Run(workerMux) }() if err := generationWorker.Recover(context.Background()); err != nil { slog.Error("恢复 AI 任务失败", "error", err) } if err := dramaParseWorker.Recover(context.Background()); err != nil { slog.Error("恢复短剧创作解析任务失败", "error", err) } if err := scriptAnalysisWorker.Recover(context.Background()); err != nil { slog.Error("恢复剧本分析任务失败", "error", err) } serverErr := make(chan error, 1) go func() { slog.Info("API 服务已启动", "address", cfg.ServerAddress) serverErr <- httpServer.ListenAndServe() }() stop := make(chan os.Signal, 1) signal.Notify(stop, syscall.SIGINT, syscall.SIGTERM) select { case sig := <-stop: slog.Info("收到停止信号", "signal", sig.String()) case err = <-serverErr: if !errors.Is(err, http.ErrServerClosed) { slog.Error("API 服务异常退出", "error", err) os.Exit(1) } case err = <-workerErr: if err != nil { slog.Error("AI 任务 Worker 异常退出", "error", err) } } ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() queueServer.Shutdown() if err := httpServer.Shutdown(ctx); err != nil { slog.Error("API 服务关闭失败", "error", err) os.Exit(1) } }