Files
2026-08-25 17:59:42 +08:00

139 lines
5.3 KiB
Go

// 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)
}
}