package repository import ( "context" "errors" "time" "fashionapi/internal/model" "gorm.io/gorm" ) // IngestRepository 爬虫入库管线专属仓储:任务队列(ingest_jobs)+ nonce 防重放 // (ingest_nces)+ 走秀正式表写入(brand_runway / brand_runway_images)。 // // 去重由 worker 按「实体键」完成:走秀 = brand_id + season_code + collection_type, // 街拍 = city + year(与 review_repository 的 SaveRunwayFromDraft / SaveStreetSnapFromDraft // 晋升时的幂等键一致),因此一个秀/街拍只会在正式表留一行,多来源图片在晋升阶段聚合。 // // 入队与领取用同一张 ingest_jobs 表,领取靠 MySQL 的 FOR UPDATE SKIP LOCKED // 实现「多 worker 安全并发」——同一条任务只会被一个 worker 拿到,其它 worker 跳过它。 type IngestRepository interface { // Enqueue 写入一条待处理任务(payload 为原始 JSON)。 Enqueue(ctx context.Context, job *model.IngestJob) error // Claim 原子领取最多 limit 条 pending 任务并置为 processing,返回这些任务。 // 用 SKIP LOCKED 保证多 worker 不抢同一条。 Claim(ctx context.Context, limit int) ([]model.IngestJob, error) // MarkDone 标记任务成功。 MarkDone(ctx context.Context, id uint32) error // MarkFailed 标记任务失败并记录错误(attempts 自增)。 MarkFailed(ctx context.Context, id uint32, errMsg string) error // ReserveNonce 写入一次性随机串;若已存在(重放)返回 ok=false。 ReserveNonce(ctx context.Context, nonce string) (ok bool, err error) // RunwayIDByEntity 按实体键(brand_id + season_code + collection_type)查是否已存在走秀正式表; // 返回 (id, found)。worker 据此跳过已晋升秀的重爬,避免重复下载。 RunwayIDByEntity(ctx context.Context, brandID uint32, seasonCode, collectionType string) (uint32, bool, error) // CreateRunwayDraft 插入走秀草稿行(status=pending),返回自增主键。 CreateRunwayDraft(ctx context.Context, d *model.BrandRunwayDraft) (uint32, error) // CreateRunwayDraftImages 批量插入草稿图片行。 CreateRunwayDraftImages(ctx context.Context, imgs []model.BrandRunwayDraftImage) error // StreetSnapIDByEntity 按实体键(city + year)查是否已存在街拍正式表;返回 (id, found)。 StreetSnapIDByEntity(ctx context.Context, city string, year uint16) (uint32, bool, error) // CreateStreetSnapDraft 插入街拍草稿行(status=pending),返回自增主键。 CreateStreetSnapDraft(ctx context.Context, d *model.StreetSnapDraft) (uint32, error) // CreateStreetSnapDraftImages 批量插入街拍草稿图片行。 CreateStreetSnapDraftImages(ctx context.Context, imgs []model.StreetSnapDraftImage) error // ListJobs 按 id 倒序分页列出入库任务(用于后台监控页)。offset/limit 控制分页。 ListJobs(ctx context.Context, offset, limit int) ([]model.IngestJob, error) // CountJobs 返回 ingest_jobs 总条数(用于分页计算总页数)。 CountJobs(ctx context.Context) (int64, error) // RetryJob 把一条 failed 任务重置回 pending,清 last_error/locked_at,等待 worker 重新处理。 RetryJob(ctx context.Context, id uint32) error // EnqueueMediaCleanup 写入一条「清理七牛孤儿图」任务(payload 为待清理 key 的 JSON)。 // 由删除图集的服务调用,把同步的七牛删除改为异步队列,避免阻塞删除请求。 EnqueueMediaCleanup(ctx context.Context, payload string) error // ScheduleRetry 失败时调用:attempts+1,未达上限则退避后重置 pending,达上限则置 failed。 // 用于临时失败(网络抖动 / 单图下载失败)的自动重试,区别于永久失败(payload 解析错等)直接 MarkFailed。 ScheduleRetry(ctx context.Context, id uint32, errMsg string) error // ListImagePHashes 返回全部已晋升图片(runway + street)的 (id, phash, kind), // 跳过 is_deleted 与 phash=0/NULL(存量未计算)。供入库时与新增图做全局汉明比对(近似去重)。 ListImagePHashes(ctx context.Context) ([]ImagePHash, error) } // ImagePHash 已晋升图片的感知哈希摘要,供入库时全局近似去重比对。 type ImagePHash struct { ID uint32 // 图片数字主键 Phash uint64 // 感知哈希(SQL 已过滤 0/NULL) Kind string // "runway" | "street"(决定 dup_of 的 hashid 类型) } type ingestRepository struct { db *gorm.DB } // NewIngestRepository 创建入库管线仓储。 func NewIngestRepository(db *gorm.DB) IngestRepository { return &ingestRepository{db: db} } func (r *ingestRepository) Enqueue(ctx context.Context, job *model.IngestJob) error { now := uint32(time.Now().Unix()) job.CreatedAt = now job.UpdatedAt = now job.Status = model.IngestStatusPending return r.db.WithContext(ctx).Create(job).Error } // Claim 在事务内 SELECT ... FOR UPDATE SKIP LOCKED 锁定 pending 行, // 立即置为 processing,再返回这些行,保证领取与状态变更原子、且不被其它 worker 重复领取。 // // 事务内先做两件事: // 1. 回收卡死的 processing 任务——worker 崩溃/被杀会留下 processing 孤儿永久卡死, // 锁定超时(IngestStuckTimeoutSec)后重置回 pending 并立即可领(next_attempt_at=now)。 // 2. 仅领取「已到重试时间」的 pending(next_attempt_at <= now),未到退避点的暂不领。 func (r *ingestRepository) Claim(ctx context.Context, limit int) ([]model.IngestJob, error) { if limit < 1 { limit = 10 } now := uint32(time.Now().Unix()) tx := r.db.WithContext(ctx).Begin() if tx.Error != nil { return nil, tx.Error } defer func() { if tx.Error != nil { tx.Rollback() } }() // 1) 回收卡死的 processing 任务(仅当时间足够大,避免服务器启动初期把刚领取的任务误回收)。 if now > model.IngestStuckTimeoutSec { if err := tx.Model(&model.IngestJob{}). Where("status = ? AND locked_at > 0 AND locked_at < ?", model.IngestStatusProcessing, now-model.IngestStuckTimeoutSec). Updates(map[string]any{ "status": model.IngestStatusPending, "locked_at": 0, "next_attempt_at": now, "updated_at": now, }).Error; err != nil { return nil, err } } // 2) 领取 pending 且已到重试时间的任务(含刚回收的 + 新入队的 + 退避已到期的)。 // 优先处理 media_cleanup(删除图集时异步清理七牛),避免被大批 crawl 任务排到后面、清理迟迟不触发。 var ids []uint32 if err := tx.Raw( "SELECT id FROM ingest_jobs WHERE status = ? AND next_attempt_at <= ? "+ "ORDER BY CASE kind WHEN ? THEN 0 ELSE 1 END, id ASC LIMIT ? FOR UPDATE SKIP LOCKED", model.IngestStatusPending, now, model.IngestKindMediaCleanup, limit, ).Scan(&ids).Error; err != nil { return nil, err } if len(ids) == 0 { tx.Commit() return nil, nil } if err := tx.Model(&model.IngestJob{}). Where("id IN ?", ids). Updates(map[string]any{ "status": model.IngestStatusProcessing, "locked_at": now, "updated_at": now, }).Error; err != nil { return nil, err } if err := tx.Commit().Error; err != nil { return nil, err } var jobs []model.IngestJob if err := r.db.WithContext(ctx).Where("id IN ?", ids).Find(&jobs).Error; err != nil { return nil, err } return jobs, nil } func (r *ingestRepository) MarkDone(ctx context.Context, id uint32) error { return r.db.WithContext(ctx). Model(&model.IngestJob{}). Where("id = ?", id). Updates(map[string]any{ "status": model.IngestStatusDone, "last_error": "", "updated_at": uint32(time.Now().Unix()), }).Error } func (r *ingestRepository) MarkFailed(ctx context.Context, id uint32, errMsg string) error { return r.db.WithContext(ctx). Model(&model.IngestJob{}). Where("id = ?", id). Updates(map[string]any{ "status": model.IngestStatusFailed, "attempts": gorm.Expr("attempts + 1"), "last_error": errMsg, "updated_at": uint32(time.Now().Unix()), }).Error } // ScheduleRetry 失败时调度自动重试:attempts+1,未达上限(IngestMaxAttempts)则按指数退避 // 重置为 pending 并写入 next_attempt_at(到点才可被 Claim 领取);达上限则置 failed,需人工处理。 // 用于临时失败(网络抖动 / 单图下载失败),区别于永久失败(payload 解析错等)直接 MarkFailed。 func (r *ingestRepository) ScheduleRetry(ctx context.Context, id uint32, errMsg string) error { now := uint32(time.Now().Unix()) var job model.IngestJob if err := r.db.WithContext(ctx).Select("attempts").Where("id = ?", id).First(&job).Error; err != nil { return err } attempts := int(job.Attempts) + 1 if attempts >= model.IngestMaxAttempts { return r.db.WithContext(ctx). Model(&model.IngestJob{}). Where("id = ?", id). Updates(map[string]any{ "status": model.IngestStatusFailed, "attempts": attempts, "last_error": errMsg, "updated_at": now, }).Error } delay := model.IngestRetryBackoff(attempts) return r.db.WithContext(ctx). Model(&model.IngestJob{}). Where("id = ?", id). Updates(map[string]any{ "status": model.IngestStatusPending, "attempts": attempts, "last_error": errMsg, "next_attempt_at": now + uint32(delay), "locked_at": 0, "updated_at": now, }).Error } // EnqueueMediaCleanup 写入一条「清理七牛孤儿图」任务(kind=media_cleanup), // payload 为待清理 key 的 JSON,由 worker 的 processMediaCleanup 按引用计数判定真孤儿后删除。 func (r *ingestRepository) EnqueueMediaCleanup(ctx context.Context, payload string) error { now := uint32(time.Now().Unix()) job := &model.IngestJob{ Kind: model.IngestKindMediaCleanup, Payload: payload, CreatedAt: now, UpdatedAt: now, Status: model.IngestStatusPending, } return r.db.WithContext(ctx).Create(job).Error } // ReserveNonce 写入一次性随机串;依赖 ingest_nces.nonce 主键唯一约束, // 重复插入触发 DuplicateEntry → 视为重放,返回 ok=false。 func (r *ingestRepository) ReserveNonce(ctx context.Context, nonce string) (bool, error) { err := r.db.WithContext(ctx).Create(&model.IngestNonce{ Nonce: nonce, CreatedAt: uint32(time.Now().Unix()), }).Error if err != nil { // 唯一键冲突 → 重放 if errors.Is(err, gorm.ErrDuplicatedKey) || isDuplicateKey(err) { return false, nil } return false, err } return true, nil } // RunwayIDByEntity 按实体键查走秀正式表是否已存在(与 SaveRunwayFromDraft 晋升键一致)。 func (r *ingestRepository) RunwayIDByEntity(ctx context.Context, brandID uint32, seasonCode, collectionType string) (uint32, bool, error) { var row struct { ID uint32 `gorm:"column:id"` } err := r.db.WithContext(ctx). Model(&model.BrandRunway{}). Select("id"). Where("brand_id = ? AND season_code = ? AND collection_type = ? AND is_deleted = 0", brandID, seasonCode, collectionType). Limit(1). Scan(&row).Error if err != nil { return 0, false, err } if row.ID == 0 { return 0, false, nil } return row.ID, true, nil } func (r *ingestRepository) CreateRunwayDraft(ctx context.Context, d *model.BrandRunwayDraft) (uint32, error) { now := uint32(time.Now().Unix()) d.CreatedAt = now d.UpdatedAt = now if d.Status == "" { d.Status = model.DraftStatusPending } if err := r.db.WithContext(ctx).Create(d).Error; err != nil { return 0, err } return d.ID, nil } func (r *ingestRepository) CreateRunwayDraftImages(ctx context.Context, imgs []model.BrandRunwayDraftImage) error { if len(imgs) == 0 { return nil } now := uint32(time.Now().Unix()) for i := range imgs { imgs[i].CreatedAt = now imgs[i].UpdatedAt = now } return r.db.WithContext(ctx).Create(&imgs).Error } // StreetSnapIDByEntity 按实体键(city + year)查街拍正式表是否已存在(与 SaveStreetSnapFromDraft 晋升键一致)。 func (r *ingestRepository) StreetSnapIDByEntity(ctx context.Context, city string, year uint16) (uint32, bool, error) { var row struct { ID uint32 `gorm:"column:id"` } err := r.db.WithContext(ctx). Model(&model.StreetSnap{}). Select("id"). Where("city = ? AND year = ? AND is_deleted = 0", city, year). Limit(1). Scan(&row).Error if err != nil { return 0, false, err } if row.ID == 0 { return 0, false, nil } return row.ID, true, nil } func (r *ingestRepository) CreateStreetSnapDraft(ctx context.Context, d *model.StreetSnapDraft) (uint32, error) { now := uint32(time.Now().Unix()) d.CreatedAt = now d.UpdatedAt = now if d.Status == "" { d.Status = model.DraftStatusPending } if err := r.db.WithContext(ctx).Create(d).Error; err != nil { return 0, err } return d.ID, nil } func (r *ingestRepository) CreateStreetSnapDraftImages(ctx context.Context, imgs []model.StreetSnapDraftImage) error { if len(imgs) == 0 { return nil } now := uint32(time.Now().Unix()) for i := range imgs { imgs[i].CreatedAt = now imgs[i].UpdatedAt = now } return r.db.WithContext(ctx).Create(&imgs).Error } // ListJobs 按 id 倒序分页列出入库任务(后台监控页用)。offset/limit 控制分页区间。 func (r *ingestRepository) ListJobs(ctx context.Context, offset, limit int) ([]model.IngestJob, error) { if limit < 1 { limit = 50 } if offset < 0 { offset = 0 } var jobs []model.IngestJob if err := r.db.WithContext(ctx).Order("id DESC").Offset(offset).Limit(limit).Find(&jobs).Error; err != nil { return nil, err } return jobs, nil } // CountJobs 返回 ingest_jobs 总条数(后台监控页分页用)。 func (r *ingestRepository) CountJobs(ctx context.Context) (int64, error) { var n int64 if err := r.db.WithContext(ctx).Model(&model.IngestJob{}).Count(&n).Error; err != nil { return 0, err } return n, nil } // RetryJob 把一条 failed 任务重置回 pending,清掉 last_error / locked_at / next_attempt_at, // 让 worker(每 3s 扫一次 pending)立即重新拉起处理。仅对 failed 生效,其它状态原样不动。 func (r *ingestRepository) RetryJob(ctx context.Context, id uint32) error { now := uint32(time.Now().Unix()) return r.db.WithContext(ctx). Model(&model.IngestJob{}). Where("id = ? AND status = ?", id, model.IngestStatusFailed). Updates(map[string]any{ "status": model.IngestStatusPending, "locked_at": 0, "last_error": "", "next_attempt_at": 0, "updated_at": now, }).Error } // ListImagePHashes 返回全部已晋升图片(runway + street)的感知哈希摘要,供入库时全局近似去重。 // 跳过 is_deleted 与 phash=0/NULL(存量未计算)。结果合并 runway + street 两类, // 用 Kind 标注类型,调用方据此把 dup_of 编码成对应 hashid 类型。 // // 注:每次入库任务都会全量拉取一次(图片量当前为千级,可接受);若后续图片量到十万级, // 可改为按 runway_id/snap_id 分批或加内存缓存 + 定时刷新,避免每 job 一次全表扫描。 func (r *ingestRepository) ListImagePHashes(ctx context.Context) ([]ImagePHash, error) { type phRow struct { ID uint32 `gorm:"column:id"` Phash uint64 `gorm:"column:phash"` } out := make([]ImagePHash, 0, 64) const whereActive = "is_deleted = 0 AND phash IS NOT NULL AND phash <> 0" var rw []phRow if err := r.db.WithContext(ctx).Model(&model.BrandRunwayImage{}). Select("id, phash").Where(whereActive).Scan(&rw).Error; err != nil { return nil, err } for _, x := range rw { out = append(out, ImagePHash{ID: x.ID, Phash: x.Phash, Kind: "runway"}) } var sn []phRow if err := r.db.WithContext(ctx).Model(&model.StreetSnapImage{}). Select("id, phash").Where(whereActive).Scan(&sn).Error; err != nil { return nil, err } for _, x := range sn { out = append(out, ImagePHash{ID: x.ID, Phash: x.Phash, Kind: "street"}) } return out, nil } // isDuplicateKey 兜底:gorm 的 ErrDuplicatedKey 在不同驱动下的封装不一定一致, // 直接命中 MySQL 1062 错误号更稳。 func isDuplicateKey(err error) bool { if err == nil { return false } msg := err.Error() return containsAny(msg, "Duplicate entry", "1062", "UNIQUE constraint failed") } func containsAny(s string, subs ...string) bool { for _, sub := range subs { if len(sub) > 0 && indexOf(s, sub) >= 0 { return true } } return false } func indexOf(s, sub string) int { for i := 0; i+len(sub) <= len(s); i++ { if s[i:i+len(sub)] == sub { return i } } return -1 }