fix(publish): 入库实体键读错误 fail-closed,复用加 rejected 守卫,去重排除软删行

This commit is contained in:
toom1996
2026-09-23 10:41:16 +08:00
parent dba615d636
commit 8f8b06def7
3 changed files with 117 additions and 4 deletions

View File

@ -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

View File

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

View File

@ -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)