Files
backend_v2/docs/superpowers/plans/2026-09-20-ingest-idempotency.md
2026-09-20 01:06:13 +08:00

1211 lines
43 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# 入库幂等实现计划(来源级去重 + 街拍实体键引入月份)
> **面向 AI 代理的工作者:** 必需子技能:使用 subagent-driven-development(推荐)或 executing-plans 逐任务实现此计划。步骤使用复选框(`- [ ]`)语法来跟踪进度。
**目标:** 让「同一篇来源文章重爬」成为幂等空操作(不下载、不写重复草稿),并修掉街拍实体键 `(city, year)` 因 `year=0` 导致的「每城永远一个专辑」缺陷。
**架构:** 契约层新增 `source` / `source_url` / `month` 三个字段;草稿表落 `(source, source_url)` 部分唯一索引作为幂等的硬保证;worker 在 `process()` 中「分派前查重、命中即 `MarkDone`」;街拍实体键扩为 `(city, year, month)`,四处落点必须同改以保证晋升与去重口径一致。
**技术栈:** Go 1.26、Gin、GORM + PostgreSQL(pgvector)、`ingest_jobs` 队列(`FOR UPDATE SKIP LOCKED`);爬虫侧独立 Go module `my-spiders`。
**规格:** `docs/superpowers/specs/2026-09-20-ingest-idempotency-design.md`
## 全局约束
- 来源幂等键固定为 `(source, source_url)`;`source` 取值 `vogue` / `theimpression`。
- 「已存在」判定**不限 `status`、不限 `is_deleted`**(决策:爬过即终态,`rejected` 也算)。
- 来源唯一索引**必须是部分索引**:`WHERE source_url <> ''`(存量行 `source_url` 为空串,全量唯一索引会因 `('','')` 冲突导致建索引失败、启动崩溃)。
- 来源唯一索引**不加** `is_deleted` 条件。
- 三个新字段 JSON 名逐字为 `source` / `source_url` / `month`,两侧(`internal/dto` 与 `spider/internal/ingest`)必须一致。
- 街拍实体键 `(city, year, month)` 共 **4 处**落点,必须同改。
- 街拍 `year` / `month` **均取抓取时刻**,不再从标题解析;`(*TheImpressionSpider).parseYear` 删除(`yearRe` 被 `vogue.go` 共用,保留)。
- 公开 API(`PublicStreetSnap` 等)**不加** `month` 字段。
- 所有新增列用 `not null default ''` / `not null default 0`,不得可空。
- 错误处理:唯一冲突一律按「跳过(`MarkDone`)」处理,**绝不**走 `failOrRetry`。
## 已完成的准备工作
本次计划**不含**存量数据清理任务:2026-09-20 已执行完毕(删除 `street_snap_drafts` 3 行 + `street_snap_draft_images` 177 行 + `ingest_jobs` 95–98,S4 孤儿 59 个对象经 `media_cleanup` 任务清理)。`brand_*` 与 job 90–94 保留未动。
## 文件结构
| 文件 | 职责 | 本计划动作 |
| --- | --- | --- |
| `internal/dto/ingest.go` | 爬虫上报载荷的后端定义(跨进程契约) | 修改:加 3 字段 |
| `internal/dto/ingest_test.go` | 契约 JSON 名锁定 | 创建 |
| `internal/model/runway_draft.go` | 走秀草稿模型 | 修改:加 source / source_url |
| `internal/model/street_snap_draft.go` | 街拍草稿模型 | 修改:加 source / source_url / month |
| `internal/model/street_snap.go` | 街拍正式表模型 | 修改:加 month |
| `internal/database/postgres.go` | 连接 + 迁移(`AutoMigrate` / `EnsureDedupSchema`) | 修改:加 2 个部分唯一索引 |
| `internal/repository/ingest_repository.go` | 入库仓储(队列 + 草稿写入 + 实体键查询) | 修改:加 `ErrDuplicateSource`、`SourceDraftExists`;`StreetSnapIDByEntity` 加 month |
| `internal/repository/review_repository.go` | 审核仓储(晋升 + 可编辑白名单) | 修改:晋升键加 month、白名单加 month |
| `internal/service/ingest_service.go` | 入库管线业务逻辑 | 修改:来源查重、填来源/月份、冲突兜底 |
| `internal/service/review_service.go` | 审核服务(字段渲染) | 修改:street 字段加 month、卡片副标题加月份 |
| `internal/repository/source_idempotency_integration_test.go` | 索引与查询的集成验证 | 创建(`//go:build integration`) |
| `internal/service/ingest_source_dedup_test.go` | 来源去重、字段透传、冲突兜底的单元验证 | 创建 |
| `internal/service/review_service_test.go` | 街拍卡片副标题的单元验证 | 创建 |
| `../spider/internal/ingest/payload.go` | 爬虫侧载荷定义 | 修改:加 3 字段 |
| `../spider/internal/ingest/payload_test.go` | 爬虫侧契约锁定 | 创建 |
| `../spider/internal/spider/vogue.go` | Vogue 爬虫 | 修改:填 source / source_url |
| `../spider/internal/spider/theimpression.go` | theImpression 街拍爬虫 | 修改:填来源与抓取时间;删 `parseYear`;抽出可测载荷构造 |
| `../spider/internal/spider/street_payload_test.go` | 街拍载荷构造的单元验证 | 创建 |
---
### 任务 1:契约三字段(两侧载荷 + JSON 契约锁定)
**文件:**
- 修改:`internal/dto/ingest.go:20-35`
- 修改:`../spider/internal/ingest/payload.go:18-33`
- 测试:创建 `internal/dto/ingest_test.go`
- 测试:创建 `../spider/internal/ingest/payload_test.go`
- [ ] **步骤 1:编写失败的后端契约测试**
创建 `internal/dto/ingest_test.go`:
```go
package dto
import (
"encoding/json"
"testing"
)
// TestRunwayIngestSourceFieldsJSON 锁定来源幂等与月份字段的 JSON 名。
// 这三者是与爬虫的跨进程契约:json tag 写错不会编译报错,
// 只会让后端永远收到空串、幂等静默失效,因此用测试钉死。
func TestRunwayIngestSourceFieldsJSON(t *testing.T) {
raw := `{"kind":"street","source":"theimpression","source_url":"https://x/y","month":8}`
var p RunwayIngest
if err := json.Unmarshal([]byte(raw), &p); err != nil {
t.Fatalf("unmarshal 失败: %v", err)
}
if p.Source != "theimpression" {
t.Fatalf("source 未解析: %q", p.Source)
}
if p.SourceURL != "https://x/y" {
t.Fatalf("source_url 未解析: %q", p.SourceURL)
}
if p.Month != 8 {
t.Fatalf("month 未解析: %d", p.Month)
}
}
```
- [ ] **步骤 2:运行测试确认失败**
运行:`go test ./internal/dto/ -run TestRunwayIngestSourceFieldsJSON -v`
预期:FAIL,编译错误 `p.Source undefined`
- [ ] **步骤 3:实现后端字段**
修改 `internal/dto/ingest.go`,在 `RunwayIngest` 结构体内、`Looks` 字段之前插入:
```go
Source string `json:"source"` // 来源站标识:vogue / theimpression(可空,来源级幂等键)
SourceURL string `json:"source_url"` // 文章详情页地址(可空;为空时不启用来源去重)
Month uint8 `json:"month"` // 抓取月份 1-12(street 用;runway 不填 = 0)
```
并把结构体上方注释的「通用字段」段落补两行:
```go
// - Source / SourceURL: 来源站与文章地址,构成来源级幂等键(同一篇文章只入库一次)
// - Month: 抓取月份(街拍实体键的一部分)
```
- [ ] **步骤 4:运行测试确认通过**
运行:`go test ./internal/dto/ -run TestRunwayIngestSourceFieldsJSON -v`
预期:PASS
- [ ] **步骤 5:编写失败的爬虫侧契约测试**
创建 `../spider/internal/ingest/payload_test.go`:
```go
package ingest
import (
"encoding/json"
"strings"
"testing"
)
// TestPayloadSourceFieldsJSON 锁定与后端的 JSON 契约:字段名必须逐字一致,
// 否则后端收到空串、来源去重静默失效。
func TestPayloadSourceFieldsJSON(t *testing.T) {
raw, err := json.Marshal(RunwayIngest{
Kind: KindStreet,
Source: "theimpression",
SourceURL: "https://x/y",
Month: 8,
})
if err != nil {
t.Fatalf("marshal 失败: %v", err)
}
for _, want := range []string{
`"source":"theimpression"`,
`"source_url":"https://x/y"`,
`"month":8`,
} {
if !strings.Contains(string(raw), want) {
t.Fatalf("payload 缺少 %s,实际: %s", want, raw)
}
}
}
```
- [ ] **步骤 6:运行测试确认失败**
运行:`cd ../spider && go test ./internal/ingest/ -run TestPayloadSourceFieldsJSON -v`
预期:FAIL,编译错误 `unknown field Source`
- [ ] **步骤 7:实现爬虫侧字段**
修改 `../spider/internal/ingest/payload.go` 的 `RunwayIngest`,在 `Looks` 之前插入:
```go
Source string `json:"source"` // 来源站标识:vogue / theimpression
SourceURL string `json:"source_url"` // 文章详情页地址(来源级幂等键)
Month uint8 `json:"month"` // 抓取月份 1-12(street 用;runway 不填 = 0)
```
- [ ] **步骤 8:运行测试确认通过**
运行:`cd ../spider && go test ./internal/ingest/ -v`
预期:PASS(含既有 `TestSignVerifyRoundtrip`)
- [ ] **步骤 9:Commit**
```bash
git add internal/dto/ingest.go internal/dto/ingest_test.go
git commit -m "feat(ingest): 上报载荷新增 source/source_url/month 三字段"
cd ../spider && git add internal/ingest/payload.go internal/ingest/payload_test.go
git commit -m "feat(ingest): 上报载荷新增 source/source_url/month 三字段"
```
---
### 任务 2:Schema(来源两列 + 街拍 month + 两个部分唯一索引)
**文件:**
- 修改:`internal/model/runway_draft.go:10-30`
- 修改:`internal/model/street_snap_draft.go:12-26`
- 修改:`internal/model/street_snap.go:9-19`
- 修改:`internal/database/postgres.go`(`EnsureDedupSchema`,在最后的 `return nil` 之前)
- 修改:`internal/repository/ingest_repository.go`(本任务一并加入 `SourceDraftExists`,因为测试要覆盖它)
- 测试:创建 `internal/repository/source_idempotency_integration_test.go`
- [ ] **步骤 1:编写失败的集成测试**
创建 `internal/repository/source_idempotency_integration_test.go`(复用同包 `dedup_integration_test.go` 已提供的 `testDB(t)` 与构建标签):
```go
//go:build integration
package repository
import (
"context"
"testing"
"fashionapi/internal/model"
)
// TestSourceIdempotencyIndex 验证部分唯一索引的两条关键语义:
// 1. 同 (source, source_url) 第二次插入必须被拒——这是来源幂等的硬保证;
// 2. source_url 为空串的行不受约束(老爬虫兼容;也是「全量唯一索引会建失败」的原因)。
func TestSourceIdempotencyIndex(t *testing.T) {
db := testDB(t)
newDraft := func(src, url string) *model.StreetSnapDraft {
return &model.StreetSnapDraft{
Title: "itest", City: "itest-city", Year: 2026, Month: 8,
Source: src, SourceURL: url, CreatedAt: 1, UpdatedAt: 1,
}
}
clean := func() { db.Where("title = ?", "itest").Delete(&model.StreetSnapDraft{}) }
clean()
t.Cleanup(clean)
if err := db.Create(newDraft("itest", "itest-url")).Error; err != nil {
t.Fatalf("首次插入应成功: %v", err)
}
if err := db.Create(newDraft("itest", "itest-url")).Error; err == nil {
t.Fatalf("同 (source, source_url) 的第二次插入必须被唯一索引拒绝")
}
if err := db.Create(newDraft("", "")).Error; err != nil {
t.Fatalf("空 source_url 的行不应受部分索引约束: %v", err)
}
if err := db.Create(newDraft("", "")).Error; err != nil {
t.Fatalf("空 source_url 的行应可多行共存: %v", err)
}
}
// TestSourceDraftExists 验证来源查重的口径:空 sourceURL 恒不命中;写入后命中;
// 且不限 status(决策:爬过即终态,rejected 也算)。
func TestSourceDraftExists(t *testing.T) {
db := testDB(t)
repo := NewIngestRepository(db)
ctx := context.Background()
clean := func() { db.Where("title = ?", "itest-exists").Delete(&model.BrandRunwayDraft{}) }
clean()
t.Cleanup(clean)
// 空 source_url:老爬虫路径,恒不命中(不得触发跳过)。
hit, err := repo.SourceDraftExists(ctx, "runway", "itest", "")
if err != nil {
t.Fatalf("SourceDraftExists 出错: %v", err)
}
if hit {
t.Fatalf("空 source_url 不应命中")
}
// 未写入前不命中。
if hit, err = repo.SourceDraftExists(ctx, "runway", "itest", "itest-exists-url"); err != nil || hit {
t.Fatalf("未写入前不应命中: hit=%v err=%v", hit, err)
}
// 写入一条 rejected 草稿:按「爬过即终态」,仍应命中。
row := &model.BrandRunwayDraft{
Title: "itest-exists", Source: "itest", SourceURL: "itest-exists-url",
Status: model.DraftStatusRejected, CreatedAt: 1, UpdatedAt: 1,
}
if err := db.Create(row).Error; err != nil {
t.Fatalf("插入草稿失败: %v", err)
}
if hit, err = repo.SourceDraftExists(ctx, "runway", "itest", "itest-exists-url"); err != nil || !hit {
t.Fatalf("rejected 草稿也应命中: hit=%v err=%v", hit, err)
}
}
```
- [ ] **步骤 2:运行测试确认失败**
运行:`go test -tags integration ./internal/repository/ -run 'TestSourceIdempotencyIndex|TestSourceDraftExists' -v`
预期:FAIL,编译错误 `Source field not found in type model.StreetSnapDraft` 与 `repo.SourceDraftExists undefined`
- [ ] **步骤 3:给三个模型加列**
`internal/model/runway_draft.go`,在 `BrandRunwayDraft` 的 `JobID` 之后插入:
```go
// 来源级幂等键:source=来源站,source_url=文章详情页地址。
// 与 uq_br_draft_source 部分唯一索引(WHERE source_url <> '')配套。
Source string `gorm:"column:source;type:varchar(32);not null;default:''" json:"source"`
SourceURL string `gorm:"column:source_url;type:varchar(512);not null;default:''" json:"source_url"`
```
`internal/model/street_snap_draft.go`,在 `StreetSnapDraft` 的 `JobID` 之后插入同样两行,并在 `Year` 之后插入:
```go
// Month 抓取月份 1-12,街拍实体键 (city, year, month) 的一部分。
Month uint8 `gorm:"column:month;not null;default:0" json:"month"`
```
`internal/model/street_snap.go`,在 `StreetSnap` 的 `Year` 之后插入:
```go
Month uint8 `gorm:"column:month;not null;default:0" json:"month"` // 街拍实体键的一部分
```
- [ ] **步骤 4:加部分唯一索引**
修改 `internal/database/postgres.go` 的 `EnsureDedupSchema`,在最后的 `return nil` 之前插入:
```go
// 来源级幂等:同一 (source, source_url) 只允许一行草稿。
// 必须是部分索引——存量草稿 source_url 为空串,全量唯一索引会因大量 ('','') 冲突
// 导致建索引失败、进而启动崩溃;该条件同时让存量数据豁免。
// 不加 is_deleted 条件:按「爬过即终态」的决策,软删也不应解锁幂等。
for _, stmt := range []string{
"CREATE UNIQUE INDEX IF NOT EXISTS uq_br_draft_source " +
"ON brand_runway_drafts (source, source_url) WHERE source_url <> ''",
"CREATE UNIQUE INDEX IF NOT EXISTS uq_ss_draft_source " +
"ON street_snap_drafts (source, source_url) WHERE source_url <> ''",
} {
if err := db.Exec(stmt).Error; err != nil {
return fmt.Errorf("create source idempotency index: %w", err)
}
}
```
- [ ] **步骤 5:加入 `SourceDraftExists`**
在 `internal/repository/ingest_repository.go` 的 `IngestRepository` 接口中(紧接 `StreetSnapIDByEntity` 之后)加入:
```go
// SourceDraftExists 按来源键 (source, source_url) 判断该文章是否已入库过。
// sourceURL 为空时直接返回 false(老爬虫兼容)。不限 status、不限 is_deleted,
// 与 uq_*_draft_source 部分唯一索引口径严格一致(决策:爬过即终态)。
SourceDraftExists(ctx context.Context, kind, source, sourceURL string) (bool, error)
```
在文件末尾加入实现,并在 import 块加入 `"fashionapi/internal/dto"`(该包已有 4 个文件 import dto,无循环依赖):
```go
// sourceDraftTables 来源幂等键所在的草稿表(按 payload kind 映射)。
var sourceDraftTables = map[string]string{
dto.IngestKindRunway: "brand_runway_drafts",
dto.IngestKindStreet: "street_snap_drafts",
}
// SourceDraftExists 按 (source, source_url) 判断该文章是否已入库过。
func (r *ingestRepository) SourceDraftExists(ctx context.Context, kind, source, sourceURL string) (bool, error) {
if sourceURL == "" {
return false, nil // 老爬虫不上送来源:不启用来源去重,保持旧行为
}
table, ok := sourceDraftTables[kind]
if !ok {
return false, nil // 未知 kind:防御性放行,交由后续分派报错
}
var n int64
if err := r.db.WithContext(ctx).Table(table).
Where("source = ? AND source_url = ?", source, sourceURL).
Count(&n).Error; err != nil {
return false, err
}
return n > 0, nil
}
```
- [ ] **步骤 6:运行测试确认通过**
运行:`go test -tags integration ./internal/repository/ -run 'TestSourceIdempotencyIndex|TestSourceDraftExists' -v`
预期:PASS(两个测试均通过)
- [ ] **步骤 7:回归 + Commit**
运行:`go vet ./... && go test ./...`
预期:vet 无输出;全部包 ok
```bash
git add internal/model/runway_draft.go internal/model/street_snap_draft.go \
internal/model/street_snap.go internal/database/postgres.go \
internal/repository/ingest_repository.go \
internal/repository/source_idempotency_integration_test.go
git commit -m "feat(ingest): 草稿表加来源幂等列与部分唯一索引,街拍表加 month 列"
```
---
### 任务 3:worker 前置查重(命中即跳过)
**文件:**
- 修改:`internal/service/ingest_service.go`(`process`,约 145-166 行)
- 测试:创建 `internal/service/ingest_source_dedup_test.go`
> 说明:`source_url == ""` 时不命中来源去重,这一条由**任务 2 的 `TestSourceDraftExists` 集成测试**覆盖(守卫在仓储层)。本任务只测服务层的跳过行为。
- [ ] **步骤 1:编写失败的单元测试**
创建 `internal/service/ingest_source_dedup_test.go`:
```go
package service
import (
"context"
"testing"
"fashionapi/internal/model"
"fashionapi/internal/repository"
)
// sourceDedupRepo 只实现来源去重用例涉及的方法;其余由内嵌接口兜底(不会被调用)。
type sourceDedupRepo struct {
repository.IngestRepository
exists bool
done []uint32
}
func (r *sourceDedupRepo) SourceDraftExists(ctx context.Context, kind, source, sourceURL string) (bool, error) {
return r.exists, nil
}
func (r *sourceDedupRepo) MarkDone(ctx context.Context, id uint32) error {
r.done = append(r.done, id)
return nil
}
// spyUploader 记录 Upload 调用次数,用于断言「命中来源去重时一张图都没上传」。
type spyUploader struct{ calls int }
func (u *spyUploader) Upload(ctx context.Context, key string, data []byte, contentType string) (string, error) {
u.calls++
return key, nil
}
func (u *spyUploader) Enabled() bool { return true }
// TestProcessSkipsOnSourceHit 来源幂等命中:任务置 done,且不下载、不写草稿。
func TestProcessSkipsOnSourceHit(t *testing.T) {
repo := &sourceDedupRepo{exists: true}
up := &spyUploader{}
svc := &IngestService{repo: repo, uploader: up}
job := model.IngestJob{
ID: 7,
Kind: model.IngestKindCrawl,
Payload: `{"kind":"street","source":"theimpression",` +
`"source_url":"https://x/y","city":"Copenhagen","images":["https://x/1.jpg"]}`,
}
svc.process(context.Background(), job)
if len(repo.done) != 1 || repo.done[0] != 7 {
t.Fatalf("命中来源去重应 MarkDone(job=7),实际 %v", repo.done)
}
if up.calls != 0 {
t.Fatalf("命中来源去重不应下载/上传任何图片,实际上传 %d 次", up.calls)
}
}
```
- [ ] **步骤 2:运行测试确认失败**
运行:`go test ./internal/service/ -run TestProcessSkipsOnSourceHit -v`
预期:FAIL —— `命中来源去重应 MarkDone(job=7),实际 []`
- [ ] **步骤 3:在 `process()` 中插入查重**
修改 `internal/service/ingest_service.go` 的 `process`,在 `p.Kind` 归一化之后、`log.Printf("... 开始处理 ...")` **之前**插入:
```go
// 来源级幂等:同一篇文章(source + source_url)爬过即终态,直接跳过,不下载、不写草稿。
// 放在 kind 归一化之后、分派之前:单点覆盖 runway / street 两条管线,
// 且早于品牌校验与实体键去重。老爬虫不上送 source_url 时不启用本去重。
if exists, err := s.repo.SourceDraftExists(ctx, p.Kind, p.Source, p.SourceURL); err == nil && exists {
log.Printf("[ingest] job=%d %s 来源去重命中(source=%s url=%s),跳过",
job.ID, p.Kind, p.Source, p.SourceURL)
_ = s.repo.MarkDone(ctx, job.ID)
return
}
```
- [ ] **步骤 4:运行测试确认通过**
运行:`go test ./internal/service/ -run TestProcessSkipsOnSourceHit -v`
预期:PASS
- [ ] **步骤 5:回归 + Commit**
运行:`go vet ./... && go test ./...`
预期:vet 无输出;全部包 ok
```bash
git add internal/service/ingest_service.go internal/service/ingest_source_dedup_test.go
git commit -m "feat(ingest): worker 处理前按来源键查重,命中即跳过"
```
---
### 任务 4:写草稿填来源/月份 + INSERT 唯一冲突当跳过
**文件:**
- 修改:`internal/repository/ingest_repository.go`(哨兵错误 + `CreateRunwayDraft` / `CreateStreetSnapDraft`)
- 修改:`internal/service/ingest_service.go`(`processRunway` / `processStreet`)
- 测试:修改 `internal/service/ingest_source_dedup_test.go`(追加两个用例)
- [ ] **步骤 1:编写失败的单元测试**
在 `internal/service/ingest_source_dedup_test.go` 末尾追加:
```go
// streetPathRepo 提供 processStreet 走到「写草稿」所需的最小实现,并捕获草稿与错误处理结果。
type streetPathRepo struct {
repository.IngestRepository
got *model.StreetSnapDraft
done []uint32
retried []uint32
failedAt []uint32
createErr error
}
func (r *streetPathRepo) SourceDraftExists(ctx context.Context, kind, source, sourceURL string) (bool, error) {
return false, nil // 前置查重未命中,模拟并发窗口
}
func (r *streetPathRepo) StreetSnapIDByEntity(ctx context.Context, city string, year uint16, month uint8) (uint32, bool, error) {
return 0, false, nil // 实体键未命中,继续下载
}
func (r *streetPathRepo) CreateStreetSnapDraft(ctx context.Context, d *model.StreetSnapDraft) (uint32, error) {
r.got = d
if r.createErr != nil {
return 0, r.createErr
}
return 42, nil
}
func (r *streetPathRepo) CreateStreetSnapDraftImages(ctx context.Context, imgs []model.StreetSnapDraftImage) error {
return nil
}
func (r *streetPathRepo) MarkDone(ctx context.Context, id uint32) error {
r.done = append(r.done, id)
return nil
}
func (r *streetPathRepo) MarkFailed(ctx context.Context, id uint32, errMsg string) error {
r.failedAt = append(r.failedAt, id)
return nil
}
func (r *streetPathRepo) ScheduleRetry(ctx context.Context, id uint32, errMsg string) error {
r.retried = append(r.retried, id)
return nil
}
// TestProcessStreetFillsSourceAndMonth 草稿须带上来源键与月份(月份是实体键的一部分,
// 漏传会导致所有街拍挤进 month=0 的桶)。
func TestProcessStreetFillsSourceAndMonth(t *testing.T) {
repo := &streetPathRepo{}
svc := &IngestService{repo: repo, uploader: &spyUploader{}}
job := model.IngestJob{
ID: 9,
Kind: model.IngestKindCrawl,
Payload: `{"kind":"street","source":"theimpression","source_url":"https://x/y",` +
`"city":"Copenhagen","year":2026,"month":8,"images":[]}`,
}
svc.process(context.Background(), job)
if repo.got == nil {
t.Fatalf("应写入街拍草稿")
}
if repo.got.Source != "theimpression" || repo.got.SourceURL != "https://x/y" {
t.Fatalf("草稿来源字段未透传: source=%q url=%q", repo.got.Source, repo.got.SourceURL)
}
if repo.got.Month != 8 || repo.got.Year != 2026 {
t.Fatalf("草稿月份/年份未透传: year=%d month=%d", repo.got.Year, repo.got.Month)
}
}
// TestProcessStreetDuplicateDraftTreatedAsSkip 并发下 INSERT 撞唯一索引时,
// 必须按「跳过」处理:MarkDone 且不调度重试(否则会白重试 3 次、堆出假故障)。
func TestProcessStreetDuplicateDraftTreatedAsSkip(t *testing.T) {
repo := &streetPathRepo{createErr: repository.ErrDuplicateSource}
svc := &IngestService{repo: repo, uploader: &spyUploader{}}
job := model.IngestJob{
ID: 10,
Kind: model.IngestKindCrawl,
Payload: `{"kind":"street","source":"theimpression","source_url":"https://x/y",` +
`"city":"Copenhagen","year":2026,"month":8,"images":[]}`,
}
svc.process(context.Background(), job)
if len(repo.done) != 1 || repo.done[0] != 10 {
t.Fatalf("唯一冲突应按跳过处理 MarkDone(10),实际 %v", repo.done)
}
if len(repo.retried) != 0 {
t.Fatalf("唯一冲突不得调度重试,实际 %v", repo.retried)
}
if len(repo.failedAt) != 0 {
t.Fatalf("唯一冲突不得置失败,实际 %v", repo.failedAt)
}
}
```
- [ ] **步骤 2:运行测试确认失败**
运行:`go test ./internal/service/ -run 'TestProcessStreetFillsSourceAndMonth|TestProcessStreetDuplicateDraftTreatedAsSkip' -v`
预期:FAIL —— `草稿来源字段未透传: source="" url=""`
- [ ] **步骤 3:加哨兵错误并在草稿写入处翻译**
修改 `internal/repository/ingest_repository.go`,在 import 之后、`IngestRepository` 接口之前加入:
```go
// ErrDuplicateSource 草稿表来源唯一键冲突:同一 (source, source_url) 已存在。
// 由 CreateRunwayDraft / CreateStreetSnapDraft 在检测到重复键时返回,
// 供服务层按「跳过」而非「失败」处理(并发窗口下的兜底)。
var ErrDuplicateSource = errors.New("duplicate source draft")
```
把 `CreateRunwayDraft` 的写入改为:
```go
if err := r.db.WithContext(ctx).Create(d).Error; err != nil {
if isDuplicateKey(err) {
return 0, ErrDuplicateSource
}
return 0, err
}
return d.ID, nil
```
`CreateStreetSnapDraft` 做同样处理。
- [ ] **步骤 4:运行测试确认错误已可识别**
运行:`go test ./internal/service/ -run 'TestProcessStreetFillsSourceAndMonth|TestProcessStreetDuplicateDraftTreatedAsSkip' -v`
预期:FAIL —— `草稿来源字段未透传: source="" url=""`(编译已通过,服务层尚未填字段)
- [ ] **步骤 5:服务层填来源/月份,并在冲突时跳过**
修改 `internal/service/ingest_service.go`。在 import 块加入 `"errors"`。
`processRunway`:在 `draft := &model.BrandRunwayDraft{` 字面量内的 `JobID: job.ID,` 之后加:
```go
Source: p.Source,
SourceURL: p.SourceURL,
```
并把其下方的错误分支改为:
```go
id, err := s.repo.CreateRunwayDraft(ctx, draft)
if err != nil {
// 唯一冲突=并发下这篇文章已由另一个任务写入草稿。按「跳过」处理,绝不能走 failOrRetry,
// 否则会按指数退避白重试 3 次、在后台堆出一批假故障。
// 此处无需 cleanupUploads:冲突行的图 URL 与本次完全相同 → sha1 内容寻址推出同一对象 key
// → 属覆盖写,不产生孤儿文件。
if errors.Is(err, repository.ErrDuplicateSource) {
log.Printf("[ingest] job=%d runway 来源并发冲突(草稿已存在),跳过", job.ID)
_ = s.repo.MarkDone(ctx, job.ID)
return
}
s.failOrRetry(ctx, job.ID, "create draft: "+err.Error())
return
}
```
`processStreet`:草稿字面量内、`JobID: job.ID,` 之后加:
```go
Source: p.Source,
SourceURL: p.SourceURL,
Month: p.Month,
```
错误分支改为:
```go
id, err := s.repo.CreateStreetSnapDraft(ctx, draft)
if err != nil {
if errors.Is(err, repository.ErrDuplicateSource) {
log.Printf("[ingest] job=%d street 来源并发冲突(草稿已存在),跳过", job.ID)
_ = s.repo.MarkDone(ctx, job.ID)
return
}
s.failOrRetry(ctx, job.ID, "create street draft: "+err.Error())
return
}
```
- [ ] **步骤 6:运行测试确认通过**
运行:`go test ./internal/service/ -v`
预期:全部 PASS
- [ ] **步骤 7:回归 + Commit**
运行:`go vet ./... && go test ./...`
预期:vet 无输出;全部包 ok
```bash
git add internal/repository/ingest_repository.go internal/service/ingest_service.go \
internal/service/ingest_source_dedup_test.go
git commit -m "feat(ingest): 草稿落来源与月份字段,唯一冲突按跳过而非失败处理"
```
---
### 任务 5:爬虫侧填充 source / source_url / year / month
**文件:**
- 修改:`../spider/internal/spider/vogue.go`(载荷构造)
- 修改:`../spider/internal/spider/theimpression.go`
- 测试:创建 `../spider/internal/spider/street_payload_test.go`
- [ ] **步骤 1:编写失败的街拍载荷测试**
创建 `../spider/internal/spider/street_payload_test.go`:
```go
package spider
import (
"testing"
"time"
)
// TestBuildStreetPayload 验证街拍载荷的来源键与抓取时间推导:
// year / month 必须同源(同取抓取时刻),否则会产生 (city, 2024, 9) 这类无意义实体键。
func TestBuildStreetPayload(t *testing.T) {
s := &TheImpressionSpider{}
now := time.Date(2026, 8, 14, 10, 30, 0, 0, time.Local)
p := s.buildStreetPayload("The Best Street Style From Copenhagen Fashion Week",
"Copenhagen", []string{"https://x/1.jpg"}, "https://theimpression.com/a/", now)
if p.Source != "theimpression" {
t.Fatalf("source 应为 theimpression,实际 %q", p.Source)
}
if p.SourceURL != "https://theimpression.com/a/" {
t.Fatalf("source_url 未填充: %q", p.SourceURL)
}
if p.Year != 2026 {
t.Fatalf("year 应取抓取年份 2026,实际 %d", p.Year)
}
if p.Month != 8 {
t.Fatalf("month 应取抓取月份 8,实际 %d", p.Month)
}
if p.City != "Copenhagen" || len(p.Images) != 1 {
t.Fatalf("city / images 未正确透传: %+v", p)
}
}
```
- [ ] **步骤 2:运行测试确认失败**
运行:`cd ../spider && go test ./internal/spider/ -run TestBuildStreetPayload -v`
预期:FAIL,编译错误 `s.buildStreetPayload undefined`
- [ ] **步骤 3:实现载荷构造并改 `getDetail`**
修改 `../spider/internal/spider/theimpression.go`:
把 `getDetail` 中的
```go
year := s.parseYear(title)
city := s.parseCity(title)
```
改为
```go
city := s.parseCity(title)
```
把上报分支的载荷构造与提交替换为:
```go
p := s.buildStreetPayload(title, city, imageURLs, url, time.Now())
if _, err := s.Ingest.Submit(context.Background(), p); err != nil {
log.Printf("[错误] 街拍入库上报失败 [%s]: %v", title, err)
} else {
log.Printf("[成功] 已上报街拍待入库: %s (共 %d 张图片)", title, len(imageURLs))
}
```
删除 `(*TheImpressionSpider).parseYear` 整个方法(已无调用方;`yearRe` 被 `vogue.go` 共用,保留)。
在文件末尾加入:
```go
// buildStreetPayload 组装街拍入库载荷。
//
// year 与 month 均取抓取时刻 now,二者必须同源:theImpression 标题实测不含年份,
// 若年份取自标题而月份取自抓取时间,会产出 (city, 2024, 9) 这类无意义实体键。
// now 作为参数传入而非内部调 time.Now(),以便单测固定时间。
func (s *TheImpressionSpider) buildStreetPayload(title, city string, images []string, sourceURL string, now time.Time) ingest.RunwayIngest {
return ingest.RunwayIngest{
Kind: ingest.KindStreet,
Source: "theimpression",
SourceURL: sourceURL,
TitleEn: title,
Year: uint16(now.Year()),
Month: uint8(now.Month()),
City: city,
Images: images,
}
}
```
并把 `getDetail` 上方那段已过时的注释替换为:
```go
// 入库分支:启用入库管线(s.Ingest != nil)时,把街拍元数据 HMAC 签名 POST 到后台
// ingest,由后台 worker 异步下载图、传 S4、写 street_snap_draft 待审;通过后才晋升
// street_snap 正式表。判重有两层:来源键 (source, source_url) 挡同篇重爬,
// 实体键 (city, year, month) 挡同城同月的重复收敛。
```
- [ ] **步骤 4:运行测试确认通过**
运行:`cd ../spider && go test ./internal/spider/ -v`
预期:PASS
- [ ] **步骤 5:给 vogue 载荷补来源字段**
修改 `../spider/internal/spider/vogue.go`,在载荷字面量 `Kind: ingest.KindRunway,` 之后插入:
```go
Source: "vogue",
SourceURL: requestURL,
```
- [ ] **步骤 6:编译并回归**
运行:`cd ../spider && gofmt -l . && go vet ./... && go test ./...`
预期:gofmt 无输出;vet 无输出;全部包 ok
- [ ] **步骤 7:Commit**
```bash
cd ../spider
git add internal/ingest/payload.go internal/ingest/payload_test.go \
internal/spider/vogue.go internal/spider/theimpression.go \
internal/spider/street_payload_test.go
git commit -m "feat(spider): 上送 source/source_url;街拍 year+month 改取抓取时间"
```
---
### 任务 6:街拍实体键引入 month(4 处落点同改)
**文件:**
- 修改:`internal/repository/ingest_repository.go`(落点 1)
- 修改:`internal/service/ingest_service.go`(落点 2)
- 修改:`internal/repository/review_repository.go`(落点 3、4)
- 测试:修改 `internal/repository/source_idempotency_integration_test.go`
- [ ] **步骤 1:编写失败的集成测试**
在 `internal/repository/source_idempotency_integration_test.go` 末尾追加:
```go
// TestStreetSnapEntityKeyIncludesMonth 街拍实体键含月份:
// 同 city + year 但 month 不同不得互判重(否则同城一年两季时装周会整季被丢弃)。
func TestStreetSnapEntityKeyIncludesMonth(t *testing.T) {
db := testDB(t)
repo := NewIngestRepository(db)
ctx := context.Background()
snap := &model.StreetSnap{Title: "itest", City: "itest-city2", Year: 2026, Month: 2, CreatedAt: 1, UpdatedAt: 1}
if err := db.Create(snap).Error; err != nil {
t.Fatalf("插入正式街拍失败: %v", err)
}
t.Cleanup(func() { db.Where("city = ?", "itest-city2").Delete(&model.StreetSnap{}) })
if _, found, err := repo.StreetSnapIDByEntity(ctx, "itest-city2", 2026, 8); err != nil || found {
t.Fatalf("不同月份不应命中: found=%v err=%v", found, err)
}
if id, found, err := repo.StreetSnapIDByEntity(ctx, "itest-city2", 2026, 2); err != nil || !found || id != snap.ID {
t.Fatalf("同月份应命中且返回该行 id=%d: id=%d found=%v err=%v", snap.ID, id, found, err)
}
}
```
- [ ] **步骤 2:运行测试确认失败**
运行:`go test -tags integration ./internal/repository/ -run TestStreetSnapEntityKeyIncludesMonth -v`
预期:FAIL,编译错误 `too many arguments in call to repo.StreetSnapIDByEntity`
- [ ] **步骤 3:落点 1 —— 实体键查询**
修改 `internal/repository/ingest_repository.go`:
接口签名与注释改为:
```go
// StreetSnapIDByEntity 按实体键(city + year + month)查是否已存在街拍正式表;返回 (id, found)。
StreetSnapIDByEntity(ctx context.Context, city string, year uint16, month uint8) (uint32, bool, error)
```
实现改为:
```go
// StreetSnapIDByEntity 按实体键(city + year + month)查街拍正式表是否已存在
// (与 SaveStreetSnapFromDraft 晋升键一致)。
func (r *ingestRepository) StreetSnapIDByEntity(ctx context.Context, city string, year uint16, month uint8) (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 month = ? AND is_deleted = 0", city, year, month).
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
}
```
- [ ] **步骤 4:落点 2 —— worker 调用点与日志**
修改 `internal/service/ingest_service.go` 的 `processStreet`:
```go
// 1) 按实体键(city + year + month)去重:正式表已存在该街拍则跳过,避免重爬重复下载。
// 只判正式表、不判 pending 草稿,保留多来源街拍图片在晋升阶段聚合。
dedupStart := time.Now()
if _, found, err := s.repo.StreetSnapIDByEntity(ctx, p.City, p.Year, p.Month); err == nil && found {
log.Printf("[ingest] job=%d street 实体去重命中(city=%s year=%d month=%d),跳过 耗时=%v",
job.ID, p.City, p.Year, p.Month, time.Since(dedupStart).Round(time.Millisecond))
_ = s.repo.MarkDone(ctx, job.ID)
return
}
```
- [ ] **步骤 5:落点 3 —— 晋升时的正式行查找**
修改 `internal/repository/review_repository.go` 的 `SaveStreetSnapFromDraft`。
把函数上方注释里的 `city+year` 改为 `city+year+month`。正式行查找改为:
```go
if eErr := tx.Where("city = ? AND year = ? AND month = ? AND is_deleted = 0", draft.City, draft.Year, draft.Month).
Limit(1).Find(&existing).Error; eErr != nil {
return eErr
}
```
`common` map 增加 `month`:
```go
common := map[string]any{
"title": draft.Title,
"year": draft.Year,
"month": draft.Month,
"city": draft.City,
"cover": draft.Cover,
"image_count": uint16(len(imgs)),
"updated_at": now,
}
```
新建分支的结构体加 `Month`:
```go
snap := &model.StreetSnap{
Title: draft.Title,
Year: draft.Year,
Month: draft.Month,
City: draft.City,
Cover: draft.Cover,
ImageCount: uint16(len(imgs)),
CreatedAt: now,
UpdatedAt: now,
}
```
- [ ] **步骤 6:落点 4 —— sibling 聚合查询(最易漏的一处)**
修改 `internal/repository/review_repository.go`:
`SaveStreetSnapFromDraft` 内的调用改为:
```go
sibImgs, err := r.streetApprovedSiblingImages(ctx, draftID, draft.City, draft.Year, draft.Month)
```
函数签名与注释改为:
```go
// streetApprovedSiblingImages 同 city+year+month 已 approved 的其他街拍草稿图(多来源聚合)。
// 月份必须与 SaveStreetSnapFromDraft 的正式行查找口径一致,
// 否则会出现「正式行按月分开了、图片却跨月混在一起」的撕裂。
func (r *reviewRepository) streetApprovedSiblingImages(ctx context.Context, excludeDraftID uint32, city string, year uint16, month uint8) ([]model.StreetSnapDraftImage, error) {
```
查询条件改为:
```go
Where("city = ? AND year = ? AND month = ? AND status = ? AND is_deleted = 0 AND id <> ?",
city, year, month, model.DraftStatusApproved, excludeDraftID).
```
- [ ] **步骤 7:运行测试确认通过**
运行:`go test -tags integration ./internal/repository/ -run TestStreetSnapEntityKeyIncludesMonth -v`
预期:PASS
- [ ] **步骤 8:回归并确认无遗漏落点**
运行:`go vet ./... && go test ./...`
预期:vet 无输出;全部包 ok(编译通过即证明 4 处签名已全部对齐)
再确认没有残留旧口径(`rg` 不可用则改用编辑器搜索):
```bash
rg "city = \? AND year = \?" internal/ --glob '*.go'
rg "StreetSnapIDByEntity" internal/ --glob '*.go'
```
预期:第一条无匹配;第二条仅剩接口定义、实现、`processStreet` 调用点三处,且调用均带 4 个实参。
- [ ] **步骤 9:Commit**
```bash
git add internal/repository/ingest_repository.go internal/repository/review_repository.go \
internal/service/ingest_service.go \
internal/repository/source_idempotency_integration_test.go
git commit -m "feat(street): 实体键扩为 city+year+month,四处落点同步"
```
---
### 任务 7:后台审核页展示并允许微调月份
**文件:**
- 修改:`internal/repository/review_repository.go:337-342`(`streetDraftEditable`)
- 修改:`internal/service/review_service.go`(`toStreetCards`、`streetModule.DraftDetail`)
- 测试:创建 `internal/service/review_service_test.go`
- [ ] **步骤 1:编写失败的单元测试**
创建 `internal/service/review_service_test.go`:
```go
package service
import (
"testing"
"fashionapi/internal/model"
)
// TestStreetCardsSubtitleHasMonth 街拍审核卡片副标题须带月份,
// 否则同城同年两季时装周在审核列表里无法区分。
func TestStreetCardsSubtitleHasMonth(t *testing.T) {
m := &streetModule{}
cards := m.toStreetCards([]model.StreetSnapDraft{
{ID: 1, Title: "t", City: "Copenhagen", Year: 2026, Month: 8},
})
if len(cards) != 1 {
t.Fatalf("期望 1 张卡片,实际 %d", len(cards))
}
if cards[0].Subtitle != "Copenhagen · 2026-08" {
t.Fatalf("副标题应含月份,实际 %q", cards[0].Subtitle)
}
}
// TestStreetSubtitleNoMonthWhenZero 月份为 0(存量行)时不显示 -00 这种噪音。
func TestStreetSubtitleNoMonthWhenZero(t *testing.T) {
m := &streetModule{}
cards := m.toStreetCards([]model.StreetSnapDraft{{ID: 2, City: "Paris", Year: 2026}})
if cards[0].Subtitle != "Paris · 2026" {
t.Fatalf("零月份应只显示年份,实际 %q", cards[0].Subtitle)
}
}
```
- [ ] **步骤 2:运行测试确认失败**
运行:`go test ./internal/service/ -run 'TestStreetCardsSubtitle|TestStreetSubtitleNoMonthWhenZero' -v`
预期:FAIL —— `副标题应含月份,实际 "Copenhagen"`
- [ ] **步骤 3:实现副标题**
修改 `internal/service/review_service.go`,把 `toStreetCards` 整体替换为:
```go
func (m *streetModule) toStreetCards(rows []model.StreetSnapDraft) []DraftCard {
cards := make([]DraftCard, 0, len(rows))
for _, d := range rows {
cards = append(cards, DraftCard{
Kind: dto.IngestKindStreet,
ID: d.ID,
Title: d.Title,
Subtitle: streetSubtitle(d.City, d.Year, d.Month),
Cover: m.img.Compose(d.Cover), // 后台列表封面看展示图
ImageCount: d.ImageCount,
Status: d.Status,
CreatedAt: d.CreatedAt,
UpdatedAt: d.UpdatedAt,
})
}
return cards
}
// streetSubtitle 组装街拍卡片副标题:城市 · YYYY-MM(月份为 0 时只到年份)。
// 月份是同城两季时装周(同年 2 月 FW / 8 月 SS)唯一的区分手段,故必须显示。
func streetSubtitle(city string, year uint16, month uint8) string {
base := city
if base == "" {
base = "—"
}
if year == 0 {
return base
}
stamp := yearStr(year)
if month > 0 {
stamp = fmt.Sprintf("%s-%02d", stamp, month)
}
return base + " · " + stamp
}
```
在 import 块加入 `"fmt"`(若尚未存在)。
- [ ] **步骤 4:运行测试确认通过**
运行:`go test ./internal/service/ -run 'TestStreetCardsSubtitle|TestStreetSubtitleNoMonthWhenZero' -v`
预期:两个测试均 PASS
- [ ] **步骤 5:把 month 加进审核详情可编辑字段**
修改 `internal/service/review_service.go` 的 `streetModule.DraftDetail`,把 `Fields` 替换为:
```go
Fields: []EditField{
{Name: "title", Label: "标题", Value: d.Title, Type: "text"},
{Name: "year", Label: "年份", Value: yearStr(d.Year), Type: "number"},
{Name: "month", Label: "月份", Value: monthStr(d.Month), Type: "number"},
{Name: "city", Label: "城市", Value: d.City, Type: "text"},
},
```
在文件末尾(`yearStr` 函数旁)加入:
```go
// monthStr 把 uint8 月份转字符串(0 显示为空)。
func monthStr(m uint8) string {
if m == 0 {
return ""
}
return strconv.FormatUint(uint64(m), 10)
}
```
- [ ] **步骤 6:放行编辑白名单**
修改 `internal/repository/review_repository.go` 的 `streetDraftEditable`:
```go
// streetDraftEditable 街拍草稿可微调字段白名单(title / year / month / city)。
var streetDraftEditable = map[string]bool{
"title": true,
"year": true,
"month": true,
"city": true,
}
```
- [ ] **步骤 7:回归 + Commit**
运行:`gofmt -l . && go vet ./... && go test ./...`
预期:gofmt 无输出;vet 无输出;全部包 ok
```bash
git add internal/service/review_service.go internal/repository/review_repository.go \
internal/service/review_service_test.go
git commit -m "feat(review): 街拍审核页展示月份并允许微调"
```
---
## 完成后的整体验证
- [ ] `cd d:/project/backend && gofmt -l . && go vet ./... && go test ./...` —— 无未格式化文件、vet 无输出、全部包 ok
- [ ] `cd d:/project/backend && go test -tags integration ./internal/repository/ -v` —— 集成测试全绿(需先起 pgvector:`docker compose -f scripts/pgvector/docker-compose.yml up -d`)
- [ ] `cd d:/project/spider && gofmt -l . && go vet ./... && go test ./...` —— 同上
- [ ] 启动后端,确认 `EnsureDedupSchema` 未报错(两个部分唯一索引创建成功)
- [ ] 跑一次 `./bin/spider.exe theimpression -n 1`,确认新增草稿的 `source='theimpression'`、`source_url` 非空、`month` 为当前月
- [ ] **再跑一次同样的命令**,确认日志出现「来源去重命中,跳过」,`ingest_jobs` 新任务 `status=done`,且 `street_snap_drafts` **行数不变**(这是本次改动的核心验收点)