This commit is contained in:
toom1996
2026-09-29 19:59:30 +08:00
parent d15d2a4701
commit 361ffc01af
28 changed files with 896 additions and 47 deletions

View File

@ -62,6 +62,7 @@ import (
"fmt"
"log"
"os"
"path/filepath"
"sort"
"strings"
"time"
@ -111,6 +112,8 @@ func main() {
runFixDup(os.Args[2:])
case "purge-rejected":
runPurgeRejected(os.Args[2:])
case "migrate":
runMigrate(os.Args[2:])
default:
usage()
os.Exit(2)
@ -131,6 +134,11 @@ func usage() {
把现存「已驳回(rejected)」且未下架的记录转删除:级联软删其全部图片、记录置
is_deleted=1(后台即消失),并把图片 key 入队 media_cleanup,由运行中的 worker
按引用计数清理 S4 孤儿文件(与后台点「删除(清空 S4)」等价,但一次性全局处理)
dbtool migrate [-config <yml] [-in <file.sql>]
执行 db/migrations 下的迁移脚本(按文件名排序、逐文件执行、幂等可重跑):
每个文件执行成功后记入 schema_migrations,重跑已记录的会跳过;库内「已存在」
视为已应用(也记入记录表后跳过),因此历史已手动 psql 过的迁移重跑不会报错中断。
-in 指定单个文件时,只执行该文件。本子命令替代 psql,免装客户端即可落库。
说明:
连接信息读取 configs/config.yml 的 database 段(可用 -config 覆盖)。
@ -405,6 +413,110 @@ func runPurgeRejected(args []string) {
log.Printf("✓ 完成:转删除记录 %d 条、图片 %d 张,已入队 media_cleanup(运行中 worker 将按引用计数清理 S4)", totalRecs, totalKeys)
}
// runMigrate 执行 db/migrations 下的迁移脚本(按文件名排序、逐文件执行、幂等可重跑)。
//
// 每个文件执行成功后记入 schema_migrations(name),重跑已记录的会直接跳过;
// 库内「已存在」(如 ALTER ADD COLUMN 报 column already exists,说明该迁移早先被
// 手动 psql 执行过)也视作已应用,记入记录表后跳过——这样历史已手动落库的迁移重跑
// 不会报错中断,新迁移则正常执行。
//
// 用现有的 execScript(pgx 简单查询协议,整段脚本在一个隐式事务里执行)跑,免装 psql。
// 默认执行 db/migrations/*.sql;-in 指定单文件时只执行该文件。
func runMigrate(args []string) {
fs := flag.NewFlagSet("migrate", flag.ExitOnError)
in := fs.String("in", "", "只执行单个迁移文件(默认执行 db/migrations 下全部 *.sql,按文件名排序)")
cfgPath := fs.String("config", "", "配置文件路径(默认 configs/config.yml)")
_ = fs.Parse(args)
db, dcfg, err := openDB(*cfgPath)
if err != nil {
log.Fatalf("✗ %v", err)
}
defer db.Close()
if err := ensureMigrationsTable(db); err != nil {
log.Fatalf("✗ 建迁移记录表失败: %v", err)
}
files, err := migrationFiles(*in)
if err != nil {
log.Fatalf("✗ %v", err)
}
if len(files) == 0 {
log.Fatalf("✗ 没有待执行迁移(检查 db/migrations 目录或 -in 路径)")
}
for _, path := range files {
name := filepath.Base(path)
applied, err := isMigrationApplied(db, name)
if err != nil {
log.Fatalf("✗ 查询迁移记录失败: %v", err)
}
if applied {
log.Printf("→ 跳过(已应用): %s", name)
continue
}
script, err := os.ReadFile(path)
if err != nil {
log.Fatalf("✗ 读迁移文件 %s 失败: %v", path, err)
}
log.Printf("→ 执行迁移: %s", name)
if err := execScript(context.Background(), db, string(script)); err != nil {
if strings.Contains(err.Error(), "already exists") {
log.Printf("→ 跳过(库内已存在,视为已应用): %s", name)
if rerr := recordMigration(db, name); rerr != nil {
log.Fatalf("✗ 记录迁移失败: %v", rerr)
}
continue
}
log.Fatalf("✗ 迁移失败 %s(已整体回滚): %v", name, err)
}
if err := recordMigration(db, name); err != nil {
log.Fatalf("✗ 记录迁移失败: %v", err)
}
log.Printf("✓ 已应用: %s", name)
}
log.Printf("✓ 迁移完成(%s)", dcfg.Addr())
}
// migrationFiles 返回待执行的迁移文件列表(按文件名排序)。-in 指定单文件时只返回它。
func migrationFiles(in string) ([]string, error) {
if in != "" {
return []string{in}, nil
}
matches, err := filepath.Glob("db/migrations/*.sql")
if err != nil {
return nil, err
}
sort.Strings(matches)
return matches, nil
}
// ensureMigrationsTable 建迁移记录表(幂等)。
func ensureMigrationsTable(db *sql.DB) error {
_, err := db.Exec(`CREATE TABLE IF NOT EXISTS public.schema_migrations (
name text PRIMARY KEY,
applied_at timestamptz NOT NULL DEFAULT now()
)`)
return err
}
// isMigrationApplied 该迁移是否已记录(已成功执行)。
func isMigrationApplied(db *sql.DB, name string) (bool, error) {
var n int
err := db.QueryRow(`SELECT count(*) FROM public.schema_migrations WHERE name = $1`, name).Scan(&n)
if err != nil {
return false, err
}
return n > 0, nil
}
// recordMigration 把迁移记入记录表(幂等,重跑不报错)。
func recordMigration(db *sql.DB, name string) error {
_, err := db.Exec(`INSERT INTO public.schema_migrations (name) VALUES ($1) ON CONFLICT (name) DO NOTHING`, name)
return err
}
// selectUint32s 执行返回单列 uint32 的查询。
func selectUint32s(db *sql.DB, q string, args ...any) ([]uint32, error) {
rows, err := db.Query(q, args...)

View File

@ -8,7 +8,7 @@ import (
"context"
"errors"
"flag"
"log"
"log/slog"
"net/http"
"os"
"os/signal"
@ -41,6 +41,19 @@ func s4BaseURL(cfgBase string, up *storage.S4Uploader) string {
return cfgBase
}
// setupLogger 根据运行模式初始化结构化日志(JSON 输出,便于集中采集)。
func setupLogger(mode string) {
level := slog.LevelDebug
switch mode {
case "release":
level = slog.LevelInfo
case "test":
level = slog.LevelWarn
}
handler := slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{Level: level})
slog.SetDefault(slog.New(handler))
}
func main() {
configPath := flag.String("config", "", "配置文件路径,默认查找 configs/config.yml")
flag.Parse()
@ -48,35 +61,46 @@ func main() {
// 1. 加载配置(yml 为主,环境变量可覆盖)
cfg, err := config.Load(*configPath)
if err != nil {
log.Fatalf("✗ 加载配置失败: %v", err)
slog.Error("加载配置失败", "error", err)
os.Exit(1)
}
if from := cfg.LoadedFrom(); from != "" {
log.Println("✓ 配置文件:", from)
} else {
log.Println("! 未找到配置文件,使用默认值与环境变量")
// 缺失的必填/推荐敏感配置:本地 TTY 交互提示补填(写回 config.local.yml),
// 容器/CI 等非交互环境缺失必填项则直接报错退出并提示用环境变量注入。
if err := config.PromptMissing(cfg); err != nil {
slog.Error("启动前校验失败", "error", err)
os.Exit(1)
}
gin.SetMode(cfg.Server.Mode)
setupLogger(cfg.Server.Mode)
if from := cfg.LoadedFrom(); from != "" {
slog.Info("配置文件", "path", from)
} else {
slog.Warn("未找到配置文件,使用默认值与环境变量")
}
// 2. 连接数据库
db, err := database.New(cfg.Database)
if err != nil {
log.Fatalf("✗ %v", err)
slog.Error("数据库连接失败", "error", err)
os.Exit(1)
}
defer func() {
if err := database.Close(db); err != nil {
log.Printf("! 关闭数据库连接失败: %v", err)
slog.Error("关闭数据库连接失败", "error", err)
}
}()
// 2.5 结构不再随启动自动迁移:全库「结构 + 索引 + 数据」统一由 cmd/dbtool 导出的
// 纯 SQL 维护(dbtool dump → psql -f)。这里只探活,不执行任何 DDL。
log.Println("✓ PostgreSQL 已连接:", cfg.Database.Addr())
slog.Info("PostgreSQL 已连接", "addr", cfg.Database.Addr())
// 3. 确保上传目录存在(静态文件服务的根目录)
uploadDir, _ := filepath.Abs(cfg.Upload.Dir)
if err := os.MkdirAll(uploadDir, 0o755); err != nil {
log.Fatalf("✗ 创建上传目录失败: %v", err)
slog.Error("创建上传目录失败", "error", err)
os.Exit(1)
}
log.Printf("✓ 静态资源: %s → %s", cfg.Upload.URLPrefix, uploadDir)
slog.Info("静态资源", "prefix", cfg.Upload.URLPrefix, "dir", uploadDir)
// 4.5 初始化公开 ID 混淆(HashID)。
// 对外接口的 id / brand_id 一律用 hashid 编码串,内部仍用数字主键;
@ -137,9 +161,9 @@ func main() {
// free 走公开 CoreIX 样式、VIP 走预签名原图,既方便换域名又能做分级(见 internal/pkg/imgurl)。
imgUp = s4Up
imgDel = s4Up
log.Printf("✓ 图片存储: 缤纷云 S4 bucket=%s base=%s(库只存 key,渲染时拼 base_url)", cfg.S4.Bucket, s4Up.PublicBaseURL())
slog.Info("图片存储(S4)", "bucket", cfg.S4.Bucket, "base", s4Up.PublicBaseURL())
} else {
log.Printf("✓ 图片存储: 本地 %s", cfg.Upload.URLPrefix)
slog.Info("图片存储(本地)", "prefix", cfg.Upload.URLPrefix)
}
// 爬虫入库服务:注入 mediaRepo(引用计数删孤儿)+ imgDel(S4/本地删除器)。
// 必须在 articleSvc / snapSvc 之前创建——图集服务删除时需调用它把「清理存储孤儿图」异步入队。
@ -155,9 +179,9 @@ func main() {
Brand: handler.NewBrandHandler(brandSvc),
StreetSnap: handler.NewStreetSnapHandler(snapSvc, favSvc),
Auth: handler.NewAuthHandler(authSvc),
Favorite: handler.NewFavoriteHandler(favSvc),
History: handler.NewHistoryHandler(histSvc),
Health: handler.NewHealthHandler(),
Favorite: handler.NewFavoriteHandler(favSvc),
History: handler.NewHistoryHandler(histSvc),
Health: handler.NewHealthHandler(),
})
// 5.5 装配管理后台(Backstage)handler,供下方独立引擎使用。
@ -172,9 +196,10 @@ func main() {
}
go func() {
log.Printf("🚀 公开服务已启动: http://localhost:%s", cfg.Server.Port)
slog.Info("公开服务已启动", "addr", "http://localhost:"+cfg.Server.Port)
if err := srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
log.Fatalf("✗ 公开服务启动失败: %v", err)
slog.Error("公开服务启动失败", "error", err)
os.Exit(1)
}
}()
@ -191,9 +216,10 @@ func main() {
}
go func() {
log.Printf("🔒 SSG 内部服务已启动: http://localhost:%s (仅构建期/回环可达)", cfg.Server.SSGPort)
slog.Info("SSG 内部服务已启动", "addr", "http://localhost:"+cfg.Server.SSGPort)
if err := srvSSG.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
log.Fatalf("✗ SSG 内部服务启动失败: %v", err)
slog.Error("SSG 内部服务启动失败", "error", err)
os.Exit(1)
}
}()
@ -216,15 +242,16 @@ func main() {
// 用独立可取消 ctx,优雅关闭时随主流程一起退出。
workerCtx, workerCancel := context.WithCancel(context.Background())
go ingestSvc.Run(workerCtx, 2, 3*time.Second)
log.Printf("⚙ 入库 worker 已启动(队列 ingest_jobs)")
slog.Info("入库 worker 已启动", "queue", "ingest_jobs")
srvBackstage := &http.Server{
Addr: ":" + cfg.Server.BackstagePort,
Handler: backstageEngine,
}
go func() {
log.Printf("🛠 管理后台已启动: http://localhost:%s/admin", cfg.Server.BackstagePort)
slog.Info("管理后台已启动", "addr", "http://localhost:"+cfg.Server.BackstagePort+"/admin")
if err := srvBackstage.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
log.Fatalf("✗ 管理后台启动失败: %v", err)
slog.Error("管理后台启动失败", "error", err)
os.Exit(1)
}
}()
@ -233,20 +260,20 @@ func main() {
quit := make(chan os.Signal, 1)
signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
<-quit
log.Println("→ 正在关闭服务 ...")
slog.Info("正在关闭服务")
ctx, cancel := context.WithTimeout(context.Background(),
time.Duration(cfg.Server.ShutdownTimeout)*time.Second)
defer cancel()
workerCancel() // 通知入库 worker 停止领取新任务
if err := srv.Shutdown(ctx); err != nil {
log.Printf("! 公开服务优雅关闭超时: %v", err)
slog.Error("公开服务优雅关闭超时", "error", err)
}
if err := srvSSG.Shutdown(ctx); err != nil {
log.Printf("! SSG 内部服务优雅关闭超时: %v", err)
slog.Error("SSG 内部服务优雅关闭超时", "error", err)
}
if err := srvBackstage.Shutdown(ctx); err != nil {
log.Printf("! 管理后台优雅关闭超时: %v", err)
slog.Error("管理后台优雅关闭超时", "error", err)
}
log.Println("✓ 服务已停止")
slog.Info("服务已停止")
}