From dba615d636bcd1e7d94735691ef705590b9e9859 Mon Sep 17 00:00:00 2001 From: toom1996 <23cm.cn@gmail.com> Date: Wed, 23 Sep 2026 10:21:51 +0800 Subject: [PATCH] =?UTF-8?q?feat(publish):=20=E7=88=AC=E8=99=AB=E5=85=A5?= =?UTF-8?q?=E5=BA=93=E7=9B=B4=E5=86=99=E6=AD=A3=E5=BC=8F=E8=A1=A8=EF=BC=8C?= =?UTF-8?q?=E9=A9=B3=E5=9B=9E=E8=A1=8C=E5=8F=AF=E5=A4=8D=E7=94=A8=E9=87=8D?= =?UTF-8?q?=E5=AE=A1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/repository/ingest_repository.go | 198 ++++++++++++------ .../ingest_single_table_integration_test.go | 131 ++++++++++++ internal/service/ingest_service.go | 131 +++++++----- internal/service/ingest_service_test.go | 23 ++ 4 files changed, 365 insertions(+), 118 deletions(-) create mode 100644 internal/repository/ingest_single_table_integration_test.go diff --git a/internal/repository/ingest_repository.go b/internal/repository/ingest_repository.go index 7758010..7451e4a 100644 --- a/internal/repository/ingest_repository.go +++ b/internal/repository/ingest_repository.go @@ -14,9 +14,9 @@ import ( // IngestRepository 爬虫入库管线专属仓储:任务队列(ingest_jobs)+ nonce 防重放 // (ingest_nonces)+ 走秀正式表写入(brand_runways / brand_runway_images)。 // -// 去重由 worker 按「实体键」完成:走秀 = brand_id + season_code + collection_type, -// 街拍 = city + year(与 review_repository 的 SaveRunwayFromDraft / SaveStreetSnapFromDraft -// 晋升时的幂等键一致),因此一个秀/街拍只会在正式表留一行,多来源图片在晋升阶段聚合。 +// 实体键:走秀 = brand_id + season_code + collection_type,街拍 = city + year。 +// 单表发布模型下入库直写正式表(status=pending),实体键命中后按既有行状态分三支: +// pending/published → 放弃本次入库;rejected → 复用该行重审。因此一个秀/街拍始终只留一行。 // // 入队与领取用同一张 ingest_jobs 表,领取靠 PostgreSQL 的 FOR UPDATE SKIP LOCKED // 实现「多 worker 安全并发」——同一条任务只会被一个 worker 拿到,其它 worker 跳过它。 @@ -32,19 +32,20 @@ type IngestRepository interface { 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 + // RunwayEntityState 按实体键(brand_id + season_code + collection_type)查正式表既有行,返回其状态。 + // worker 据此分三支:pending/published → 直接放弃本次入库;rejected → 复用该行重审;未命中 → 新建。 + RunwayEntityState(ctx context.Context, brandID uint32, seasonCode, collectionType string) (id uint32, status string, found bool, err error) + // CreateRunwayWithImages 新建走秀正式行(status=pending)并写入图片,返回记录主键(事务内完成)。 + CreateRunwayWithImages(ctx context.Context, rw *model.BrandRunway, imgs []model.BrandRunwayImage) (uint32, error) + // ReuseRejectedRunway 复用一条 rejected 走秀:覆盖内容字段、软删旧图、写入新图、置回 pending(事务内完成)。 + // 驳回因此不是永久黑名单:重爬同一实体即重新送审,且始终「一个实体一行」。 + ReuseRejectedRunway(ctx context.Context, id uint32, rw *model.BrandRunway, imgs []model.BrandRunwayImage) error + // StreetSnapEntityState 按实体键(city + year)查正式表既有行,返回其状态。 + StreetSnapEntityState(ctx context.Context, city string, year uint16) (id uint32, status string, found bool, err error) + // CreateStreetSnapWithImages 新建街拍正式行(status=pending)并写入图片(事务内完成)。 + CreateStreetSnapWithImages(ctx context.Context, snap *model.StreetSnap, imgs []model.StreetSnapImage) (uint32, error) + // ReuseRejectedStreetSnap 复用一条 rejected 街拍记录(事务内完成)。 + ReuseRejectedStreetSnap(ctx context.Context, id uint32, snap *model.StreetSnap, imgs []model.StreetSnapImage) error // FindNearDuplicateImage 在给定图片表中按 dHash 汉明距离检索近重复,返回命中行 id(仅取最近一条)。 FindNearDuplicateImage(ctx context.Context, tables []string, phashBits string, threshold int) (dupID uint32, found bool, err error) // ListJobs 按 id 倒序分页列出入库任务(用于后台监控页)。offset/limit 控制分页。 @@ -239,94 +240,167 @@ func (r *ingestRepository) ReserveNonce(ctx context.Context, nonce string) (bool return true, nil } -// RunwayIDByEntity 按实体键查走秀正式表是否已存在(与 SaveRunwayFromDraft 晋升键一致)。 -func (r *ingestRepository) RunwayIDByEntity(ctx context.Context, brandID uint32, seasonCode, collectionType string) (uint32, bool, error) { +// ---- 入库直写正式表(单表发布模型)---- + +func (r *ingestRepository) RunwayEntityState(ctx context.Context, brandID uint32, seasonCode, collectionType string) (uint32, string, bool, error) { var row struct { - ID uint32 `gorm:"column:id"` + ID uint32 `gorm:"column:id"` + Status string `gorm:"column:status"` } err := r.db.WithContext(ctx). Model(&model.BrandRunway{}). - Select("id"). + Select("id, status"). 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 + return 0, "", false, err } if row.ID == 0 { - return 0, false, nil + return 0, "", false, nil } - return row.ID, true, nil + return row.ID, row.Status, true, nil } -func (r *ingestRepository) CreateRunwayDraft(ctx context.Context, d *model.BrandRunwayDraft) (uint32, error) { +func (r *ingestRepository) CreateRunwayWithImages(ctx context.Context, rw *model.BrandRunway, imgs []model.BrandRunwayImage) (uint32, error) { now := uint32(time.Now().Unix()) - d.CreatedAt = now - d.UpdatedAt = now - if d.Status == "" { - d.Status = model.DraftStatusPending + rw.CreatedAt, rw.UpdatedAt = now, now + if rw.Status == "" { + rw.Status = model.StatusPending } - if err := r.db.WithContext(ctx).Create(d).Error; err != nil { + err := r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + if cErr := tx.Create(rw).Error; cErr != nil { + return cErr + } + if len(imgs) == 0 { + return nil + } + for i := range imgs { + imgs[i].RunwayID = rw.ID + imgs[i].CreatedAt, imgs[i].UpdatedAt = now, now + } + return tx.Create(&imgs).Error + }) + if err != nil { return 0, err } - return d.ID, nil + return rw.ID, nil } -func (r *ingestRepository) CreateRunwayDraftImages(ctx context.Context, imgs []model.BrandRunwayDraftImage) error { - if len(imgs) == 0 { - return nil - } +func (r *ingestRepository) ReuseRejectedRunway(ctx context.Context, id uint32, rw *model.BrandRunway, imgs []model.BrandRunwayImage) error { 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 + return r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + if uErr := tx.Model(&model.BrandRunway{}).Where("id = ?", id).Updates(map[string]any{ + "title_en": rw.TitleEn, + "title_cn": rw.TitleCn, + "description_en": rw.DescriptionEn, + "description_cn": rw.DescriptionCn, + "year": rw.Year, + "season": rw.Season, + "collection_type": rw.CollectionType, + "season_code": rw.SeasonCode, + "cover": rw.Cover, + "image_count": rw.ImageCount, + "job_id": rw.JobID, + "status": model.StatusPending, + "reviewer": "", + "reject_reason": "", + "updated_at": now, + }).Error; uErr != nil { + return uErr + } + if dErr := tx.Model(&model.BrandRunwayImage{}). + Where("runway_id = ? AND is_deleted = 0", id). + Updates(map[string]any{"is_deleted": 1, "updated_at": now}).Error; dErr != nil { + return dErr + } + if len(imgs) == 0 { + return nil + } + for i := range imgs { + imgs[i].RunwayID = id + imgs[i].CreatedAt, imgs[i].UpdatedAt = now, now + } + return tx.Create(&imgs).Error + }) } -// StreetSnapIDByEntity 按实体键(city + year)查街拍正式表是否已存在(与 SaveStreetSnapFromDraft 晋升键一致)。 -func (r *ingestRepository) StreetSnapIDByEntity(ctx context.Context, city string, year uint16) (uint32, bool, error) { +func (r *ingestRepository) StreetSnapEntityState(ctx context.Context, city string, year uint16) (uint32, string, bool, error) { var row struct { - ID uint32 `gorm:"column:id"` + ID uint32 `gorm:"column:id"` + Status string `gorm:"column:status"` } err := r.db.WithContext(ctx). Model(&model.StreetSnap{}). - Select("id"). + Select("id, status"). Where("city = ? AND year = ? AND is_deleted = 0", city, year). Limit(1). Scan(&row).Error if err != nil { - return 0, false, err + return 0, "", false, err } if row.ID == 0 { - return 0, false, nil + return 0, "", false, nil } - return row.ID, true, nil + return row.ID, row.Status, true, nil } -func (r *ingestRepository) CreateStreetSnapDraft(ctx context.Context, d *model.StreetSnapDraft) (uint32, error) { +func (r *ingestRepository) CreateStreetSnapWithImages(ctx context.Context, snap *model.StreetSnap, imgs []model.StreetSnapImage) (uint32, error) { now := uint32(time.Now().Unix()) - d.CreatedAt = now - d.UpdatedAt = now - if d.Status == "" { - d.Status = model.DraftStatusPending + snap.CreatedAt, snap.UpdatedAt = now, now + if snap.Status == "" { + snap.Status = model.StatusPending } - if err := r.db.WithContext(ctx).Create(d).Error; err != nil { + err := r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + if cErr := tx.Create(snap).Error; cErr != nil { + return cErr + } + if len(imgs) == 0 { + return nil + } + for i := range imgs { + imgs[i].SnapID = snap.ID + imgs[i].CreatedAt, imgs[i].UpdatedAt = now, now + } + return tx.Create(&imgs).Error + }) + if err != nil { return 0, err } - return d.ID, nil + return snap.ID, nil } -func (r *ingestRepository) CreateStreetSnapDraftImages(ctx context.Context, imgs []model.StreetSnapDraftImage) error { - if len(imgs) == 0 { - return nil - } +func (r *ingestRepository) ReuseRejectedStreetSnap(ctx context.Context, id uint32, snap *model.StreetSnap, imgs []model.StreetSnapImage) error { 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 + return r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + if uErr := tx.Model(&model.StreetSnap{}).Where("id = ?", id).Updates(map[string]any{ + "title": snap.Title, + "year": snap.Year, + "city": snap.City, + "cover": snap.Cover, + "image_count": snap.ImageCount, + "job_id": snap.JobID, + "status": model.StatusPending, + "reviewer": "", + "reject_reason": "", + "updated_at": now, + }).Error; uErr != nil { + return uErr + } + if dErr := tx.Model(&model.StreetSnapImage{}). + Where("snap_id = ? AND is_deleted = 0", id). + Updates(map[string]any{"is_deleted": 1, "updated_at": now}).Error; dErr != nil { + return dErr + } + if len(imgs) == 0 { + return nil + } + for i := range imgs { + imgs[i].SnapID = id + imgs[i].CreatedAt, imgs[i].UpdatedAt = now, now + } + return tx.Create(&imgs).Error + }) } // ListJobs 按 id 倒序分页列出入库任务(后台监控页用)。offset/limit 控制分页区间。 diff --git a/internal/repository/ingest_single_table_integration_test.go b/internal/repository/ingest_single_table_integration_test.go new file mode 100644 index 0000000..0c2dd0b --- /dev/null +++ b/internal/repository/ingest_single_table_integration_test.go @@ -0,0 +1,131 @@ +//go:build integration + +// 集成测试:入库直写正式表的三条分支(新建 / 放弃 / 复用驳回行)。 +package repository + +import ( + "context" + "testing" + + "fashionapi/internal/model" +) + +func newRunwayForIngest(jobID uint32, brandID uint32, seasonCode, collectionType string, cover string) *model.BrandRunway { + return &model.BrandRunway{ + JobID: jobID, BrandID: brandID, TitleEn: "ingest-test", + SeasonCode: seasonCode, CollectionType: collectionType, + Cover: cover, ImageCount: 1, Status: model.StatusPending, + } +} + +func newRunwayImages(cover string) []model.BrandRunwayImage { + return []model.BrandRunwayImage{{BrandID: 1, Image: cover, Name: "Look 1", SortOrder: 1, LookIndex: 1}} +} + +// TestIngestCreatesPendingRecord 新建分支:实体键不存在 → 写入正式表且 status=pending,图片落到正式图片表。 +func TestIngestCreatesPendingRecord(t *testing.T) { + db := testDB(t) + applyMigration(t, db, "2026-09-22-01-single-table-publish.sql") + repo := NewIngestRepository(db) + ctx := context.Background() + + const season = "SS97" + _, _, found, err := repo.RunwayEntityState(ctx, 1, season, "rtw") + if err != nil { + t.Fatalf("RunwayEntityState 出错: %v", err) + } + if found { + t.Fatalf("测试前置:实体键 %s 不应已存在", season) + } + + id, err := repo.CreateRunwayWithImages(ctx, newRunwayForIngest(101, 1, season, "rtw", "ingest-ss97.jpg"), newRunwayImages("ingest-ss97.jpg")) + if err != nil { + t.Fatalf("CreateRunwayWithImages 出错: %v", err) + } + t.Cleanup(func() { + db.Exec("DELETE FROM brand_runway_images WHERE runway_id = ?", id) + db.Exec("DELETE FROM brand_runways WHERE id = ?", id) + }) + + var rw model.BrandRunway + if err := db.Where("id = ?", id).First(&rw).Error; err != nil { + t.Fatalf("读回记录失败: %v", err) + } + if rw.Status != model.StatusPending { + t.Fatalf("新入库记录应为 pending,实际 %s", rw.Status) + } + + gotID, status, found, err := repo.RunwayEntityState(ctx, 1, season, "rtw") + if err != nil { + t.Fatalf("RunwayEntityState 出错: %v", err) + } + if !found || gotID != id || status != model.StatusPending { + t.Fatalf("实体键应命中刚建的行,实际 found=%v id=%d status=%s", found, gotID, status) + } + + publicCount := int64(0) + db.Table("public_brand_runways").Where("id = ?", id).Count(&publicCount) + if publicCount != 0 { + t.Fatalf("pending 记录不应出现在公开视图,实际 %d", publicCount) + } +} + +// TestIngestReusesRejectedRecord 复用分支:实体键命中 rejected → 覆盖内容、软删旧图、写新图、置回 pending。 +func TestIngestReusesRejectedRecord(t *testing.T) { + db := testDB(t) + applyMigration(t, db, "2026-09-22-01-single-table-publish.sql") + repo := NewIngestRepository(db) + ctx := context.Background() + + const season = "FW96" + oldID, err := repo.CreateRunwayWithImages(ctx, newRunwayForIngest(90, 1, season, "rtw", "old.jpg"), newRunwayImages("old.jpg")) + if err != nil { + t.Fatalf("预置记录失败: %v", err) + } + t.Cleanup(func() { + db.Exec("DELETE FROM brand_runway_images WHERE runway_id = ?", oldID) + db.Exec("DELETE FROM brand_runways WHERE id = ?", oldID) + }) + if err := db.Model(&model.BrandRunway{}).Where("id = ?", oldID). + Updates(map[string]any{"status": model.StatusRejected, "reviewer": "admin", "reject_reason": "图片缺失"}).Error; err != nil { + t.Fatalf("置为已驳回失败: %v", err) + } + + _, status, found, err := repo.RunwayEntityState(ctx, 1, season, "rtw") + if err != nil || !found || status != model.StatusRejected { + t.Fatalf("前置:应命中 rejected 行,实际 found=%v status=%s err=%v", found, status, err) + } + + newRW := newRunwayForIngest(91, 1, season, "rtw", "new.jpg") + newRW.TitleEn = "ingest-test-reused" + if err := repo.ReuseRejectedRunway(ctx, oldID, newRW, newRunwayImages("new.jpg")); err != nil { + t.Fatalf("ReuseRejectedRunway 出错: %v", err) + } + + var rw model.BrandRunway + if err := db.Where("id = ?", oldID).First(&rw).Error; err != nil { + t.Fatalf("读回记录失败: %v", err) + } + if rw.Status != model.StatusPending { + t.Fatalf("复用后应回到 pending,实际 %s", rw.Status) + } + if rw.TitleEn != "ingest-test-reused" || rw.Cover != "new.jpg" { + t.Fatalf("内容应被覆盖,实际 title=%s cover=%s", rw.TitleEn, rw.Cover) + } + if rw.Reviewer != "" || rw.RejectReason != "" { + t.Fatalf("复用应清空审核痕迹,实际 reviewer=%s reason=%s", rw.Reviewer, rw.RejectReason) + } + + var alive []model.BrandRunwayImage + if err := db.Where("runway_id = ? AND is_deleted = 0", oldID).Find(&alive).Error; err != nil { + t.Fatalf("读图片失败: %v", err) + } + if len(alive) != 1 || alive[0].Image != "new.jpg" { + t.Fatalf("旧图应被软删、只留新图,实际 %d 张", len(alive)) + } + var oldAlive int64 + db.Model(&model.BrandRunwayImage{}).Where("runway_id = ? AND image = ? AND is_deleted = 0", oldID, "old.jpg").Count(&oldAlive) + if oldAlive != 0 { + t.Fatalf("旧图应被软删") + } +} diff --git a/internal/service/ingest_service.go b/internal/service/ingest_service.go index 45ce68a..40c0d33 100644 --- a/internal/service/ingest_service.go +++ b/internal/service/ingest_service.go @@ -229,7 +229,7 @@ func (s *IngestService) failOrRetry(ctx context.Context, id uint32, errMsg strin } } -// processRunway 走秀入库:品牌校验 → 去重 → 补季节码 → 下载图 → 写 brand_runway_drafts。 +// processRunway 走秀入库:品牌校验 → 实体键分三支(新建 / 放弃 / 复用驳回)→ 下载图 → 直写 brand_runways。 func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob, p dto.RunwayIngest) { // 1) 品牌必须存在(爬虫负责先建/复用品牌) brandID, err := hashid.Decode(p.BrandUID) @@ -242,38 +242,35 @@ func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob, return } - // 2) 按实体键(brand_id + season_code + collection_type)去重:正式表已存在该秀则跳过, - // 避免重爬重复下载。注意:只判正式表、不判 pending 草稿——否则多来源(Vogue + theImpression) - // 爬同一场秀时第二个来源会被误判重复而丢弃,破坏晋升阶段的图片聚合。 + // 2) 实体键查重并决定分支(单表模型) seasonCode := season.Derive(p.Year, p.CollectionType, p.Season) - dedupStart := time.Now() - if _, found, err := s.repo.RunwayIDByEntity(ctx, brandID, seasonCode, p.CollectionType); err == nil && found { - log.Printf("[ingest] job=%d runway 实体去重命中(season=%s type=%s),跳过 耗时=%v", - job.ID, seasonCode, p.CollectionType, time.Since(dedupStart).Round(time.Millisecond)) + reuseID, existStatus, exist, _ := s.repo.RunwayEntityState(ctx, brandID, seasonCode, p.CollectionType) + action := decideIngestAction(exist, existStatus) + if action == ingestSkip { + log.Printf("[ingest] job=%d runway 实体已存在(status=%s),跳过", job.ID, existStatus) _ = s.repo.MarkDone(ctx, job.ID) return } - log.Printf("[ingest] job=%d runway 实体去重 耗时=%v(未命中,继续下载)", job.ID, time.Since(dedupStart).Round(time.Millisecond)) // 3) season_code 已在去重前补齐(vogue.go 历史漏填的 bug,统一在此兜底) // 4) 下载图片并上传到存储(结构化 Looks 优先:主图+细节图分组;否则回退 Images 全部视为主图)。 // 内容哈希(sha1)key 保证重爬不产生孤儿文件:失败回滚删本批 key 即可。 var cover string - var draftImages []model.BrandRunwayDraftImage + var images []model.BrandRunwayImage var keys []string var imgFailed bool var imageCount uint16 var timing *fetchTiming if len(p.Looks) > 0 { - cover, draftImages, keys, imgFailed, timing = s.fetchLookImages(ctx, p.Looks, "runway") + cover, images, keys, imgFailed, timing = s.fetchLookImages(ctx, p.Looks, "runway") imageCount = uint16(len(p.Looks)) } else { var fimgs []fetchedImage cover, fimgs, keys, imgFailed, timing = s.fetchImages(ctx, p.Images, "runway", dto.IngestKindRunway) imageCount = uint16(len(fimgs)) for i, fi := range fimgs { - draftImages = append(draftImages, model.BrandRunwayDraftImage{ + images = append(images, model.BrandRunwayImage{ Image: fi.url, Name: fmt.Sprintf("Look %d", i+1), SortOrder: uint32(i + 1), @@ -294,9 +291,9 @@ func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob, return } - // 5) 写草稿表(status=pending),等待后台审核通过后再晋升正式表 + // 5) 直写正式表(status=pending),等待后台审核通过(审核只改状态,不重建图片) writeStart := time.Now() - draft := &model.BrandRunwayDraft{ + rw := &model.BrandRunway{ JobID: job.ID, BrandID: brandID, TitleEn: p.TitleEn, @@ -309,26 +306,28 @@ func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob, SeasonCode: seasonCode, Cover: cover, ImageCount: imageCount, - Status: model.DraftStatusPending, + Status: model.StatusPending, } - id, err := s.repo.CreateRunwayDraft(ctx, draft) - if err != nil { - s.failOrRetry(ctx, job.ID, "create draft: "+err.Error()) + for i := range images { + images[i].BrandID = brandID + } + var saveErr error + if action == ingestReuseRejected { + saveErr = s.repo.ReuseRejectedRunway(ctx, reuseID, rw, images) + } else { + _, saveErr = s.repo.CreateRunwayWithImages(ctx, rw, images) + } + if saveErr != nil { + s.failOrRetry(ctx, job.ID, "save runway: "+saveErr.Error()) return } - for i := range draftImages { - draftImages[i].DraftID = id - } - if err := s.repo.CreateRunwayDraftImages(ctx, draftImages); err != nil { - s.failOrRetry(ctx, job.ID, "create draft images: "+err.Error()) - return - } - log.Printf("[ingest] job=%d runway 写草稿 耗时=%v 图片行=%d", job.ID, time.Since(writeStart).Round(time.Millisecond), len(draftImages)) + log.Printf("[ingest] job=%d runway 写正式表 耗时=%v 图片行=%d 复用驳回=%v", + job.ID, time.Since(writeStart).Round(time.Millisecond), len(images), action == ingestReuseRejected) _ = s.repo.MarkDone(ctx, job.ID) } -// fetchLookImages 按 Looks 结构下载主图+细节图并上传到存储,返回可直接落草稿的 -// BrandRunwayDraftImage 行(带 look_index / is_detail 分组)。任意一张下载/上传失败即把 +// fetchLookImages 按 Looks 结构下载主图+细节图并上传到存储,返回可直接落正式图片表的 +// BrandRunwayImage 行(带 look_index / is_detail 分组)。任意一张下载/上传失败即把 // failed 置 true,调用方据此把整条任务判失败并回滚本批已上传的 key,符合「单图失败=整任务失败」策略。 // cover 取首个成功下载的主图;image_count(主图数)由调用方按 len(Looks) 计,不在此返回。 // 每张图入库前做去重:命中 dHash 近重复则仍入库但标记留痕。 @@ -336,11 +335,11 @@ func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob, // 实现:先把全部「主图 + 细节图」按原始顺序摊平成任务列表,用 concurrentFetch 并发完成 // 「下载 → phash → 上传」这段 IO 密集操作;随后再按下标顺序串行去重与组装。 // 这样既拿到并发收益,又保证 cover / sort_order / 批次内去重与原串行实现一致。 -func (s *IngestService) fetchLookImages(ctx context.Context, looks []dto.RunwayLook, prefix string) (string, []model.BrandRunwayDraftImage, []string, bool, *fetchTiming) { +func (s *IngestService) fetchLookImages(ctx context.Context, looks []dto.RunwayLook, prefix string) (string, []model.BrandRunwayImage, []string, bool, *fetchTiming) { start := time.Now() timing := &fetchTiming{} cover := "" - rows := make([]model.BrandRunwayDraftImage, 0) + rows := make([]model.BrandRunwayImage, 0) keys := make([]string, 0) seen := make(map[string]bool) failed := false @@ -392,7 +391,7 @@ func (s *IngestService) fetchLookImages(ctx context.Context, looks []dto.RunwayL name = fmt.Sprintf("Look %d — Detail %d", task.lookIdx, task.detailNo) isDetail = 1 } - rows = append(rows, model.BrandRunwayDraftImage{ + rows = append(rows, model.BrandRunwayImage{ Image: r.url, Name: name, SortOrder: uint32(order), @@ -408,18 +407,16 @@ func (s *IngestService) fetchLookImages(ctx context.Context, looks []dto.RunwayL return cover, rows, keys, failed, timing } -// processStreet 街拍入库:去重 → 下载图 → 写 street_snap_draft(无品牌)。 +// processStreet 街拍入库:实体键分三支(新建 / 放弃 / 复用驳回)→ 下载图 → 直写 street_snaps(无品牌)。 func (s *IngestService) processStreet(ctx context.Context, job model.IngestJob, p dto.RunwayIngest) { - // 1) 按实体键(city + year)去重:正式表已存在该街拍则跳过,避免重爬重复下载。 - // 只判正式表、不判 pending 草稿,保留多来源街拍图片在晋升阶段聚合。 - dedupStart := time.Now() - if _, found, err := s.repo.StreetSnapIDByEntity(ctx, p.City, p.Year); err == nil && found { - log.Printf("[ingest] job=%d street 实体去重命中(city=%s year=%d),跳过 耗时=%v", - job.ID, p.City, p.Year, time.Since(dedupStart).Round(time.Millisecond)) + // 1) 实体键查重并决定分支(单表模型,同 runway) + reuseID, existStatus, exist, _ := s.repo.StreetSnapEntityState(ctx, p.City, p.Year) + action := decideIngestAction(exist, existStatus) + if action == ingestSkip { + log.Printf("[ingest] job=%d street 实体已存在(status=%s),跳过", job.ID, existStatus) _ = s.repo.MarkDone(ctx, job.ID) return } - log.Printf("[ingest] job=%d street 实体去重 耗时=%v(未命中,继续下载)", job.ID, time.Since(dedupStart).Round(time.Millisecond)) // 2) 下载图片并上传到存储(S4优先,失败兜底本地)。 cover, fimgs, keys, imgFailed, timing := s.fetchImages(ctx, p.Images, "street", dto.IngestKindStreet) @@ -432,26 +429,20 @@ func (s *IngestService) processStreet(ctx context.Context, job model.IngestJob, return } - // 3) 写草稿表(status=pending),等待后台审核通过后再晋升 street_snaps 正式表 + // 3) 直写正式表(status=pending) writeStart := time.Now() - draft := &model.StreetSnapDraft{ + snap := &model.StreetSnap{ JobID: job.ID, Title: p.TitleEn, // 街拍单标题,爬虫优先填 title_en Year: p.Year, City: p.City, Cover: cover, ImageCount: uint16(len(fimgs)), - Status: model.DraftStatusPending, + Status: model.StatusPending, } - id, err := s.repo.CreateStreetSnapDraft(ctx, draft) - if err != nil { - s.failOrRetry(ctx, job.ID, "create street draft: "+err.Error()) - return - } - rows := make([]model.StreetSnapDraftImage, 0, len(fimgs)) + rows := make([]model.StreetSnapImage, 0, len(fimgs)) for i, fi := range fimgs { - rows = append(rows, model.StreetSnapDraftImage{ - DraftID: id, + rows = append(rows, model.StreetSnapImage{ Image: fi.url, Name: fmt.Sprintf("Look %d", i+1), SortOrder: uint32(i + 1), @@ -460,12 +451,18 @@ func (s *IngestService) processStreet(ctx context.Context, job model.IngestJob, DupOf: strconv.FormatUint(uint64(fi.dupID), 10), }) } - - if err := s.repo.CreateStreetSnapDraftImages(ctx, rows); err != nil { - s.failOrRetry(ctx, job.ID, "create street draft images: "+err.Error()) + var saveErr error + if action == ingestReuseRejected { + saveErr = s.repo.ReuseRejectedStreetSnap(ctx, reuseID, snap, rows) + } else { + _, saveErr = s.repo.CreateStreetSnapWithImages(ctx, snap, rows) + } + if saveErr != nil { + s.failOrRetry(ctx, job.ID, "save street snap: "+saveErr.Error()) return } - log.Printf("[ingest] job=%d street 写草稿 耗时=%v 图片行=%d", job.ID, time.Since(writeStart).Round(time.Millisecond), len(rows)) + log.Printf("[ingest] job=%d street 写正式表 耗时=%v 图片行=%d 复用驳回=%v", + job.ID, time.Since(writeStart).Round(time.Millisecond), len(rows), action == ingestReuseRejected) _ = s.repo.MarkDone(ctx, job.ID) } @@ -606,9 +603,9 @@ func (s *IngestService) dedupImage(ctx context.Context, kind, phashBits string, var tables []string switch kind { case dto.IngestKindRunway: - tables = []string{"brand_runway_draft_images", "brand_runway_images"} + tables = []string{"brand_runway_images"} case dto.IngestKindStreet: - tables = []string{"street_snap_draft_images", "street_snap_images"} + tables = []string{"street_snap_images"} default: return 0, 0 } @@ -742,3 +739,25 @@ func (s *IngestService) downloadAndUpload(ctx context.Context, u, prefix string) } return res } + +// ingestActionKind 入库分支决策结果(单表发布模型)。 +type ingestActionKind int + +const ( + ingestCreate ingestActionKind = iota // 实体键未命中 → 新建记录 + ingestSkip // 命中 pending/published → 放弃本次入库 + ingestReuseRejected // 命中 rejected → 复用该行重审 +) + +// decideIngestAction 按实体键命中情况与既有行状态决定入库分支。 +// 抽成纯函数,是为了让「驳回不是永久黑名单」这条规则有独立的、不需要数据库的测试。 +func decideIngestAction(found bool, status string) ingestActionKind { + switch { + case !found: + return ingestCreate + case status == model.StatusRejected: + return ingestReuseRejected + default: + return ingestSkip + } +} diff --git a/internal/service/ingest_service_test.go b/internal/service/ingest_service_test.go index ccadcba..b1157d3 100644 --- a/internal/service/ingest_service_test.go +++ b/internal/service/ingest_service_test.go @@ -84,3 +84,26 @@ func TestIngestRetryBackoff(t *testing.T) { } } } + +// TestDecideIngestAction 入库三分支决策:未命中→新建;命中 pending/published→放弃;命中 rejected→复用重审。 +// 这是「驳回不是永久黑名单」这条规则的直接证据。 +func TestDecideIngestAction(t *testing.T) { + cases := []struct { + name string + found bool + status string + want ingestActionKind + }{ + {"未命中→新建", false, "", ingestCreate}, + {"命中 pending→放弃", true, model.StatusPending, ingestSkip}, + {"命中 published→放弃", true, model.StatusPublished, ingestSkip}, + {"命中 rejected→复用重审", true, model.StatusRejected, ingestReuseRejected}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + if got := decideIngestAction(c.found, c.status); got != c.want { + t.Fatalf("decideIngestAction(%v, %q) = %v, want %v", c.found, c.status, got, c.want) + } + }) + } +}