From 8f8b06def7bc718155f9044bb6ad8caa69c9c432 Mon Sep 17 00:00:00 2001 From: toom1996 <23cm.cn@gmail.com> Date: Wed, 23 Sep 2026 10:41:16 +0800 Subject: [PATCH] =?UTF-8?q?fix(publish):=20=E5=85=A5=E5=BA=93=E5=AE=9E?= =?UTF-8?q?=E4=BD=93=E9=94=AE=E8=AF=BB=E9=94=99=E8=AF=AF=20fail-closed?= =?UTF-8?q?=EF=BC=8C=E5=A4=8D=E7=94=A8=E5=8A=A0=20rejected=20=E5=AE=88?= =?UTF-8?q?=E5=8D=AB=EF=BC=8C=E5=8E=BB=E9=87=8D=E6=8E=92=E9=99=A4=E8=BD=AF?= =?UTF-8?q?=E5=88=A0=E8=A1=8C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/repository/ingest_repository.go | 22 ++++- .../ingest_single_table_integration_test.go | 84 +++++++++++++++++++ internal/service/ingest_service.go | 15 +++- 3 files changed, 117 insertions(+), 4 deletions(-) diff --git a/internal/repository/ingest_repository.go b/internal/repository/ingest_repository.go index 7451e4a..703dbb0 100644 --- a/internal/repository/ingest_repository.go +++ b/internal/repository/ingest_repository.go @@ -3,6 +3,7 @@ package repository import ( "context" "errors" + "fmt" "math" "time" @@ -290,7 +291,11 @@ func (r *ingestRepository) CreateRunwayWithImages(ctx context.Context, rw *model func (r *ingestRepository) ReuseRejectedRunway(ctx context.Context, id uint32, rw *model.BrandRunway, imgs []model.BrandRunwayImage) error { now := uint32(time.Now().Unix()) return r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { - if uErr := tx.Model(&model.BrandRunway{}).Where("id = ?", id).Updates(map[string]any{ + // WHERE 带 status=rejected 守卫:状态是读阶段拿到的,而写库发生在整批图片下载/上传之后 + // (可能数十秒窗口)。若无条件覆盖,期间已被改成 published 的行会被置回 pending, + // 等于把已发布内容从公开视图上撤下。 + upd := tx.Model(&model.BrandRunway{}).Where("id = ? AND status = ?", id, model.StatusRejected) + if uErr := upd.Updates(map[string]any{ "title_en": rw.TitleEn, "title_cn": rw.TitleCn, "description_en": rw.DescriptionEn, @@ -309,6 +314,9 @@ func (r *ingestRepository) ReuseRejectedRunway(ctx context.Context, id uint32, r }).Error; uErr != nil { return uErr } + if upd.RowsAffected == 0 { + return fmt.Errorf("runway %d 已不是 rejected 状态,复用前置条件失效", id) + } 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 { @@ -373,7 +381,9 @@ func (r *ingestRepository) CreateStreetSnapWithImages(ctx context.Context, snap func (r *ingestRepository) ReuseRejectedStreetSnap(ctx context.Context, id uint32, snap *model.StreetSnap, imgs []model.StreetSnapImage) error { now := uint32(time.Now().Unix()) return r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { - if uErr := tx.Model(&model.StreetSnap{}).Where("id = ?", id).Updates(map[string]any{ + // 同 runway:带 status=rejected 守卫 + 行数校验,避免把期间已发布的记录撤下公开视图。 + upd := tx.Model(&model.StreetSnap{}).Where("id = ? AND status = ?", id, model.StatusRejected) + if uErr := upd.Updates(map[string]any{ "title": snap.Title, "year": snap.Year, "city": snap.City, @@ -387,6 +397,9 @@ func (r *ingestRepository) ReuseRejectedStreetSnap(ctx context.Context, id uint3 }).Error; uErr != nil { return uErr } + if upd.RowsAffected == 0 { + return fmt.Errorf("street snap %d 已不是 rejected 状态,复用前置条件失效", id) + } 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 { @@ -448,6 +461,10 @@ func (r *ingestRepository) RetryJob(ctx context.Context, id uint32) error { // FindNearDuplicateImage 在给定图片表中按 dHash 汉明距离检索近重复,取距离 ≤ threshold 的最近一条。 // 算子用 pgvector 的 L2(<->);因 phash 是 {0,1}^64 向量,L2² == 汉明距离,故 L2 阈值 = sqrt(threshold)。 // phashBits 为 vector(64) 二进制向量串;NULL 的 phash 不参与比较。 +// +// 只比对未软删的行(is_deleted = 0):复用驳回行时去重先于软删旧图执行, +// 否则重爬到的同一张图会命中「即将被软删的旧行」,留痕 dup_of 指向一条公开不可见的记录。 +// 本方法被 runway / street 两条入库路径共用,过滤对两者语义一致。 func (r *ingestRepository) FindNearDuplicateImage(ctx context.Context, tables []string, phashBits string, threshold int) (uint32, bool, error) { // 汉明阈值转 L2 阈值:phash 为 {0,1}^64 向量,L2² == 汉明距离,故 L2 阈值 = sqrt(汉明阈值)。 l2Limit := math.Sqrt(float64(threshold)) @@ -459,6 +476,7 @@ func (r *ingestRepository) FindNearDuplicateImage(ctx context.Context, tables [] err := r.db.WithContext(ctx).Table(t). Select("id, (phash <-> ?::vector) AS dist", phashBits). Where("(phash <-> ?::vector) <= ?", phashBits, l2Limit). + Where("is_deleted = 0"). Order("dist ASC"). Limit(1). Scan(&row).Error diff --git a/internal/repository/ingest_single_table_integration_test.go b/internal/repository/ingest_single_table_integration_test.go index 0c2dd0b..9faadd1 100644 --- a/internal/repository/ingest_single_table_integration_test.go +++ b/internal/repository/ingest_single_table_integration_test.go @@ -5,9 +5,12 @@ package repository import ( "context" + "database/sql" "testing" + "time" "fashionapi/internal/model" + "fashionapi/internal/pkg/phash" ) func newRunwayForIngest(jobID uint32, brandID uint32, seasonCode, collectionType string, cover string) *model.BrandRunway { @@ -129,3 +132,84 @@ func TestIngestReusesRejectedRecord(t *testing.T) { t.Fatalf("旧图应被软删") } } + +// TestIngestReuseRejectedRefusesPublished 复用守卫:实体键状态是读阶段拿到的,而写库发生在整批 +// 图片下载/上传之后(可能数十秒窗口)。若该行在这段窗口内已被改成 published,复用必须失败, +// 绝不能把它置回 pending——那等于把已发布内容从公开视图上撤下。 +func TestIngestReuseRejectedRefusesPublished(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 = "SS95" + id, err := repo.CreateRunwayWithImages(ctx, newRunwayForIngest(80, 1, season, "rtw", "pub.jpg"), newRunwayImages("pub.jpg")) + if err != nil { + t.Fatalf("预置记录失败: %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) + }) + // 模拟「读阶段之后、写库之前」该行被审核通过。 + if err := db.Model(&model.BrandRunway{}).Where("id = ?", id). + Update("status", model.StatusPublished).Error; err != nil { + t.Fatalf("置为已发布失败: %v", err) + } + + newRW := newRunwayForIngest(81, 1, season, "rtw", "pub2.jpg") + newRW.TitleEn = "ingest-test-should-not-apply" + if err := repo.ReuseRejectedRunway(ctx, id, newRW, newRunwayImages("pub2.jpg")); err == nil { + t.Fatalf("对已发布行复用应报错,实际返回 nil") + } + + var rw model.BrandRunway + if err := db.Where("id = ?", id).First(&rw).Error; err != nil { + t.Fatalf("读回记录失败: %v", err) + } + if rw.Status != model.StatusPublished { + t.Fatalf("已发布行不得被复用改回 pending,实际 %s", rw.Status) + } + if rw.TitleEn != "ingest-test" || rw.Cover != "pub.jpg" { + t.Fatalf("复用失败时不应覆盖内容,实际 title=%s cover=%s", rw.TitleEn, rw.Cover) + } + + publicCount := int64(0) + db.Table("public_brand_runways").Where("id = ?", id).Count(&publicCount) + if publicCount != 1 { + t.Fatalf("已发布行应仍在公开视图,实际 %d", publicCount) + } +} + +// TestDedupIgnoresSoftDeletedImage 去重不得命中已软删的行:复用驳回行时去重先于软删旧图执行, +// 若把「即将被软删的旧图」算进比对,重爬到的同一张图会被打上 dup_of=<旧行 id>, +// 留痕指向一条随后公开不可见的记录。 +// +// 用与探测向量完全相同的 phash(距离 0)插入软删行,保证修复前它必然是「最近一条」被返回, +// 因此本断言在修复前确定性失败、修复后确定性通过,与库中其它数据无关。 +func TestDedupIgnoresSoftDeletedImage(t *testing.T) { + db := testDB(t) + applyMigration(t, db, "2026-09-22-01-single-table-publish.sql") + repo := NewIngestRepository(db) + ctx := context.Background() + + bits := phash.ToVectorBits(^uint64(0)) + now := uint32(time.Now().Unix()) + row := model.BrandRunwayImage{ + RunwayID: 1, BrandID: 1, Image: "ingest/soft-deleted.jpg", Name: "look", + SortOrder: 1, IsDeleted: 1, CreatedAt: now, UpdatedAt: now, + Phash: sql.NullString{String: bits, Valid: true}, + } + if err := db.Create(&row).Error; err != nil { + t.Fatalf("插入软删图片失败: %v", err) + } + t.Cleanup(func() { db.Where("id = ?", row.ID).Delete(&model.BrandRunwayImage{}) }) + + dupID, found, err := repo.FindNearDuplicateImage(ctx, []string{"brand_runway_images"}, bits, phash.DefaultThreshold) + if err != nil { + t.Fatalf("FindNearDuplicateImage 出错: %v", err) + } + if found && dupID == row.ID { + t.Fatalf("已软删的行(id=%d)不应参与近重复比对", row.ID) + } +} diff --git a/internal/service/ingest_service.go b/internal/service/ingest_service.go index 40c0d33..f9fdd05 100644 --- a/internal/service/ingest_service.go +++ b/internal/service/ingest_service.go @@ -244,7 +244,13 @@ func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob, // 2) 实体键查重并决定分支(单表模型) seasonCode := season.Derive(p.Year, p.CollectionType, p.Season) - reuseID, existStatus, exist, _ := s.repo.RunwayEntityState(ctx, brandID, seasonCode, p.CollectionType) + reuseID, existStatus, exist, err := s.repo.RunwayEntityState(ctx, brandID, seasonCode, p.CollectionType) + if err != nil { + // 读失败绝不能当成「未命中」:实体键上没有唯一约束兜底,误判为新建会给同一实体留两行, + // 破坏「一个实体一行」的不变式。fail-closed + 退避重试。 + s.failOrRetry(ctx, job.ID, "entity lookup: "+err.Error()) + return + } action := decideIngestAction(exist, existStatus) if action == ingestSkip { log.Printf("[ingest] job=%d runway 实体已存在(status=%s),跳过", job.ID, existStatus) @@ -410,7 +416,12 @@ func (s *IngestService) fetchLookImages(ctx context.Context, looks []dto.RunwayL // processStreet 街拍入库:实体键分三支(新建 / 放弃 / 复用驳回)→ 下载图 → 直写 street_snaps(无品牌)。 func (s *IngestService) processStreet(ctx context.Context, job model.IngestJob, p dto.RunwayIngest) { // 1) 实体键查重并决定分支(单表模型,同 runway) - reuseID, existStatus, exist, _ := s.repo.StreetSnapEntityState(ctx, p.City, p.Year) + reuseID, existStatus, exist, err := s.repo.StreetSnapEntityState(ctx, p.City, p.Year) + if err != nil { + // 同 runway:读失败不得当「未命中」,否则会给同一实体留两行。fail-closed + 退避重试。 + s.failOrRetry(ctx, job.ID, "entity lookup: "+err.Error()) + return + } action := decideIngestAction(exist, existStatus) if action == ingestSkip { log.Printf("[ingest] job=%d street 实体已存在(status=%s),跳过", job.ID, existStatus)