package service import ( "context" "crypto/sha1" "database/sql" "encoding/hex" "encoding/json" "fmt" "io" "log" "net/http" "path" "strconv" "time" "fashionapi/internal/dto" "fashionapi/internal/model" "fashionapi/internal/pkg/hashid" "fashionapi/internal/pkg/phash" "fashionapi/internal/pkg/season" "fashionapi/internal/pkg/storage" "fashionapi/internal/repository" ) // IngestService 爬虫入库管线的业务逻辑层:Submit 入队,worker 异步处理。 // // 队列复用一张 ingest_jobs 表,靠 kind 列分流两类任务: // - crawl:爬虫上报的走秀/街拍入库(下载图 → 写草稿) // - media_cleanup:删除图集时异步清理七牛孤儿图(引用计数归零才真删) type IngestService struct { repo repository.IngestRepository brandRepo repository.BrandRepository media repository.MediaRepository // 跨表图片引用计数(media_cleanup 删孤儿用) uploader storage.Uploader // 主上传器(七牛启用时为七牛,否则本地) del storage.Deleter // 删除器(七牛或本地,media_cleanup 真删用) local *storage.LocalUploader // 七牛失败时的兜底落地 httpClient *http.Client } // NewIngestService 创建入库服务。 // // uploader: 主上传器(七牛或本地) // del: 删除器(七牛或本地),用于 media_cleanup 真删七牛/本地孤儿文件 // local: 本地兜底上传器(七牛上传失败时回退,避免图片完全丢失) func NewIngestService(repo repository.IngestRepository, brandRepo repository.BrandRepository, media repository.MediaRepository, uploader storage.Uploader, del storage.Deleter, local *storage.LocalUploader) *IngestService { return &IngestService{ repo: repo, brandRepo: brandRepo, media: media, uploader: uploader, del: del, local: local, httpClient: &http.Client{Timeout: 30 * time.Second}, } } // Submit 把一条上报写入队列(payload 原样存 JSON),立即返回任务 id。 // 真正处理由 worker 异步完成,因此本方法本身很快、不阻塞爬虫。 func (s *IngestService) Submit(ctx context.Context, payload dto.RunwayIngest) (uint32, error) { if payload.Kind == "" { payload.Kind = dto.IngestKindRunway } raw, err := json.Marshal(payload) if err != nil { return 0, NewError(http.StatusBadRequest, "payload invalid: "+err.Error()) } // runway 必须带品牌;street 不关联品牌,brand_uid 可选 if payload.Kind == dto.IngestKindRunway && payload.BrandUID == "" { return 0, NewError(http.StatusBadRequest, "brand_uid required") } job := &model.IngestJob{Payload: string(raw), Kind: model.IngestKindCrawl} if err := s.repo.Enqueue(ctx, job); err != nil { return 0, internalErr(err.Error()) } return job.ID, nil } // ListJobs 分页列出入库任务(后台监控页用)。返回本页任务、总条数与错误。 func (s *IngestService) ListJobs(ctx context.Context, offset, limit int) ([]model.IngestJob, int64, error) { jobs, err := s.repo.ListJobs(ctx, offset, limit) if err != nil { return nil, 0, err } total, err := s.repo.CountJobs(ctx) if err != nil { return nil, 0, err } return jobs, total, nil } // RetryJob 把一条 failed 任务重新入队,等待 worker 重新处理。 func (s *IngestService) RetryJob(ctx context.Context, id uint32) error { return s.repo.RetryJob(ctx, id) } // Run 启动 worker 循环:定时领取 pending 任务并处理,直到 ctx 取消。 // 多实例安全:Claim 用 SKIP LOCKED,互不抢同一条。 func (s *IngestService) Run(ctx context.Context, batch int, interval time.Duration) { if batch < 1 { batch = 2 } if interval <= 0 { interval = 3 * time.Second } ticker := time.NewTicker(interval) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: s.drain(ctx, batch) } } } // drain 一次性领取并处理一批任务。 func (s *IngestService) drain(ctx context.Context, batch int) { jobs, err := s.repo.Claim(ctx, batch) if err != nil { return } for i := range jobs { s.process(ctx, jobs[i]) } } // process 处理单条任务:先按 job.Kind(DB 列)分流到 crawl / media_cleanup 两条管线。 // 旧爬虫任务(迁移前)kind 为空,回落 crawl 走秀/街拍入库管线。 func (s *IngestService) process(ctx context.Context, job model.IngestJob) { switch job.Kind { case model.IngestKindMediaCleanup: s.processMediaCleanup(ctx, job) return case "", model.IngestKindCrawl: // 走秀/街拍入库管线,按 payload 内 Kind 再分流 runway / street。 default: // 未知 job.Kind 绝不静默处理,标记失败避免脏数据。 _ = s.repo.MarkFailed(ctx, job.ID, "unknown job kind: "+job.Kind) return } var p dto.RunwayIngest if err := json.Unmarshal([]byte(job.Payload), &p); err != nil { _ = s.repo.MarkFailed(ctx, job.ID, "payload parse: "+err.Error()) return } if p.Kind == "" { p.Kind = dto.IngestKindRunway } switch p.Kind { case dto.IngestKindRunway: s.processRunway(ctx, job, p) case dto.IngestKindStreet: s.processStreet(ctx, job, p) default: // 未知 kind 绝不静默当成走秀处理;所有爬虫数据都必须落到已注册的审核模块, // 否则标记任务失败,避免出现「未审核就入库」的脏数据。 _ = s.repo.MarkFailed(ctx, job.ID, "unknown ingest kind: "+p.Kind) } } // MediaCleanupPayload 是「清理七牛孤儿图」任务的 payload:待清理的图片 key 列表。 type MediaCleanupPayload struct { Keys []string `json:"keys"` } // EnqueueMediaCleanup 把一批待清理的七牛 key 异步入队;真正删除由 worker 的 // processMediaCleanup 按引用计数判定,仅当 key 在所有图集/草稿表中引用归零才真删 // (内容寻址共享 key 不会被误删)。 func (s *IngestService) EnqueueMediaCleanup(ctx context.Context, keys []string) error { if len(keys) == 0 { return nil } raw, err := json.Marshal(MediaCleanupPayload{Keys: keys}) if err != nil { return err } return s.repo.EnqueueMediaCleanup(ctx, string(raw)) } // processMediaCleanup 清理七牛孤儿图任务:解析 key 列表,按跨表引用计数删除真孤儿。 // 删除/计数失败一律保守跳过(宁可留文件),因此本任务几乎总是成功置 done。 func (s *IngestService) processMediaCleanup(ctx context.Context, job model.IngestJob) { if s.media == nil || s.del == nil { _ = s.repo.MarkFailed(ctx, job.ID, "media cleanup not configured") return } var p MediaCleanupPayload if err := json.Unmarshal([]byte(job.Payload), &p); err != nil { _ = s.repo.MarkFailed(ctx, job.ID, "media cleanup payload parse: "+err.Error()) return } if len(p.Keys) == 0 { _ = s.repo.MarkDone(ctx, job.ID) return } purgeOrphanImages(ctx, s.media, s.del, p.Keys) _ = s.repo.MarkDone(ctx, job.ID) } // failOrRetry 把临时失败(如网络抖动 / 单图下载失败)按指数退避自动重试: // attempts 未达上限则重置为 pending 并推迟 next_attempt_at,worker 到点再领; // 达上限才置 failed,等待后台手动「重试」。永久失败(payload 解析错 / 未知 kind)直接 // MarkFailed,不进重试。 func (s *IngestService) failOrRetry(ctx context.Context, id uint32, errMsg string) { if err := s.repo.ScheduleRetry(ctx, id, errMsg); err != nil { // 调度失败兜底硬失败,避免任务卡在 processing 无人收。 _ = s.repo.MarkFailed(ctx, id, errMsg) } } // processRunway 走秀入库:品牌校验 → 去重 → 补季节码 → 下载图 → 写 brand_runway_draft。 func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob, p dto.RunwayIngest) { // 1) 品牌必须存在(爬虫负责先建/复用品牌) brandID, err := hashid.Decode(p.BrandUID) if err != nil { _ = s.repo.MarkFailed(ctx, job.ID, "unknown brand_uid") return } if _, err := s.brandRepo.FindByID(ctx, brandID); err != nil { _ = s.repo.MarkFailed(ctx, job.ID, "brand not found") return } // 2) 按实体键(brand_id + season_code + collection_type)去重:正式表已存在该秀则跳过, // 避免重爬重复下载。注意:只判正式表、不判 pending 草稿——否则多来源(Vogue + theImpression) // 爬同一场秀时第二个来源会被误判重复而丢弃,破坏晋升阶段的图片聚合。 seasonCode := season.Derive(p.Year, p.CollectionType, p.Season) if _, found, err := s.repo.RunwayIDByEntity(ctx, brandID, seasonCode, p.CollectionType); err == nil && found { _ = s.repo.MarkDone(ctx, job.ID) return } // 3) season_code 已在去重前补齐(vogue.go 历史漏填的 bug,统一在此兜底) // 4) 下载图片并上传到存储(结构化 Looks 优先:主图+细节图分组;否则回退 Images 全部视为主图)。 // 内容哈希(sha1)key 保证重爬不产生孤儿文件:失败回滚删本批 key 即可。 var cover string var draftImages []model.BrandRunwayDraftImage var keys []string var imgFailed bool var imageCount uint16 if len(p.Looks) > 0 { cover, draftImages, keys, imgFailed = s.fetchLookImages(ctx, p.Looks, "runway") imageCount = uint16(len(p.Looks)) } else { var fimgs []fetchedImage cover, fimgs, keys, imgFailed = s.fetchImages(ctx, p.Images, "runway", dto.IngestKindRunway) imageCount = uint16(len(fimgs)) for i, fi := range fimgs { draftImages = append(draftImages, model.BrandRunwayDraftImage{ Image: fi.url, Name: fmt.Sprintf("Look %d", i+1), SortOrder: uint32(i + 1), LookIndex: uint32(i + 1), IsDetail: 0, ContentSha1: fi.sha1, Phash: sqlNull(fi.phash), IsDuplicate: fi.isDup, DupOf: strconv.FormatUint(uint64(fi.dupID), 10), }) } } if imgFailed { // 单图失败=整任务失败:先回滚本批已上传的图(S4 + 本地兜底),避免孤儿文件永远堆在存储里, // 然后按指数退避自动重试,达上限才置 failed 等后台手动重试。 s.cleanupUploads(ctx, keys) s.failOrRetry(ctx, job.ID, "image download failed") return } // 5) 写草稿表(status=pending),等待后台审核通过后再晋升正式表 draft := &model.BrandRunwayDraft{ JobID: job.ID, BrandID: brandID, TitleEn: p.TitleEn, TitleCn: p.TitleCn, DescriptionEn: p.DescriptionEn, DescriptionCn: p.DescriptionCn, Year: p.Year, Season: p.Season, CollectionType: p.CollectionType, SeasonCode: seasonCode, Cover: cover, ImageCount: imageCount, Status: model.DraftStatusPending, } id, err := s.repo.CreateRunwayDraft(ctx, draft) if err != nil { s.failOrRetry(ctx, job.ID, "create draft: "+err.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 } _ = s.repo.MarkDone(ctx, job.ID) } // fetchLookImages 按 Looks 结构下载主图+细节图并上传到存储,返回可直接落草稿的 // BrandRunwayDraftImage 行(带 look_index / is_detail 分组)。任意一张下载/上传失败即把 // failed 置 true,调用方据此把整条任务判失败并回滚本批已上传的 key,符合「单图失败=整任务失败」策略。 // cover 取首个成功下载的主图;image_count(主图数)由调用方按 len(Looks) 计,不在此返回。 // 每张图入库前做去重:content_sha1 已存在则整行跳过(精确重复);命中 dHash 近重复则仍入库但标记留痕。 func (s *IngestService) fetchLookImages(ctx context.Context, looks []dto.RunwayLook, prefix string) (string, []model.BrandRunwayDraftImage, []string, bool) { cover := "" rows := make([]model.BrandRunwayDraftImage, 0) keys := make([]string, 0) seen := make(map[string]bool) failed := false order := 0 for li, look := range looks { lookIdx := li + 1 if look.Main != "" { url, key, sha1h, ph, err := s.downloadOne(ctx, look.Main, order, prefix) if err != nil { failed = true } else { keys = append(keys, key) if seen[sha1h] { continue } seen[sha1h] = true skip, dupID, isDup := s.dedupImage(ctx, dto.IngestKindRunway, sha1h, ph) if skip { continue } order++ if cover == "" { cover = url } rows = append(rows, model.BrandRunwayDraftImage{ Image: url, Name: fmt.Sprintf("Look %d", lookIdx), SortOrder: uint32(order), LookIndex: uint32(lookIdx), IsDetail: 0, ContentSha1: sha1h, Phash: sqlNull(ph), IsDuplicate: isDup, DupOf: strconv.FormatUint(uint64(dupID), 10), }) } } for di, d := range look.Details { if d == "" { continue } url, key, sha1h, ph, err := s.downloadOne(ctx, d, order, prefix) if err != nil { failed = true continue } keys = append(keys, key) if seen[sha1h] { continue } seen[sha1h] = true skip, dupID, isDup := s.dedupImage(ctx, dto.IngestKindRunway, sha1h, ph) if skip { continue } order++ rows = append(rows, model.BrandRunwayDraftImage{ Image: url, Name: fmt.Sprintf("Look %d — Detail %d", lookIdx, di+1), SortOrder: uint32(order), LookIndex: uint32(lookIdx), IsDetail: 1, ContentSha1: sha1h, Phash: sqlNull(ph), IsDuplicate: isDup, DupOf: strconv.FormatUint(uint64(dupID), 10), }) } } return cover, rows, keys, failed } // processStreet 街拍入库:去重 → 下载图 → 写 street_snap_draft(无品牌)。 func (s *IngestService) processStreet(ctx context.Context, job model.IngestJob, p dto.RunwayIngest) { // 1) 按实体键(city + year)去重:正式表已存在该街拍则跳过,避免重爬重复下载。 // 只判正式表、不判 pending 草稿,保留多来源街拍图片在晋升阶段聚合。 if _, found, err := s.repo.StreetSnapIDByEntity(ctx, p.City, p.Year); err == nil && found { _ = s.repo.MarkDone(ctx, job.ID) return } // 2) 下载图片并上传到存储(七牛优先,失败兜底本地)。 cover, fimgs, keys, imgFailed := s.fetchImages(ctx, p.Images, "street", dto.IngestKindStreet) if len(p.Images) > 0 && imgFailed { // 单图失败=整任务失败:先回滚本批已上传的图(七牛 + 本地兜底),避免孤儿文件永远堆在存储里, // 然后按指数退避自动重试,达上限才置 failed 等后台手动重试。 s.cleanupUploads(ctx, keys) s.failOrRetry(ctx, job.ID, "image download failed") return } // 3) 写草稿表(status=pending),等待后台审核通过后再晋升 street_snap 正式表 draft := &model.StreetSnapDraft{ JobID: job.ID, Title: p.TitleEn, // 街拍单标题,爬虫优先填 title_en Year: p.Year, City: p.City, Cover: cover, ImageCount: uint16(len(fimgs)), Status: model.DraftStatusPending, } 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)) for i, fi := range fimgs { rows = append(rows, model.StreetSnapDraftImage{ DraftID: id, Image: fi.url, Name: fmt.Sprintf("Look %d", i+1), SortOrder: uint32(i + 1), ContentSha1: fi.sha1, Phash: sqlNull(fi.phash), IsDuplicate: fi.isDup, 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()) return } _ = s.repo.MarkDone(ctx, job.ID) } // fetchImages 下载图片并上传到存储,返回 (cover 地址, 全部图片地址, 已成功上传对象的 key 列表, 是否有任意一张失败)。 // prefix 为七牛 key 前缀(runway/ 或 street/)。只要任意一张下载/上传失败,failed 即置 true, // 调用方据此把整条任务判为失败(不再写草稿),并拿 keys 回滚本批已上传的对象,符合「单图失败=整任务失败」策略。 // fetchImages 下载图片并上传到存储,返回 (cover 地址, 已下载图结构, 已成功上传对象的 key 列表, 是否有任意一张失败)。 // prefix 为七牛 key 前缀(runway/ 或 street/)。kind 用于选择去重比对表。 // 只要任意一张下载/上传失败,failed 即置 true,调用方据此把整条任务判为失败并回滚本批已上传的对象, // 符合「单图失败=整任务失败」策略。每张图入库前做去重(精确跳过 + 近重复标记)。 func (s *IngestService) fetchImages(ctx context.Context, urls []string, prefix, kind string) (string, []fetchedImage, []string, bool) { cover := "" out := make([]fetchedImage, 0, len(urls)) keys := make([]string, 0, len(urls)) seen := make(map[string]bool) failed := false for _, u := range urls { url, key, sha1h, ph, err := s.downloadOne(ctx, u, 0, prefix) if err != nil { failed = true continue } keys = append(keys, key) if seen[sha1h] { continue } seen[sha1h] = true skip, dupID, isDup := s.dedupImage(ctx, kind, sha1h, ph) if skip { continue } if cover == "" { cover = url } out = append(out, fetchedImage{ url: url, key: key, sha1: sha1h, phash: ph, dupID: dupID, isDup: isDup, }) } return cover, out, keys, failed } // cleanupUploads 删除一批本批次成功上传的对象(七牛 + 本地兜底),用于任务失败回滚: // 「单图失败=整任务失败」时,前面已成功上传的图若放任不管就成了孤儿,永远堆在存储里。 // 对两种存储都尝试删除(任一不存在即按幂等成功处理);删除失败仅告警,不阻断任务置失败。 // 注意:key 为 sha1 内容寻址,若恰与其他已晋升图集共享同一内容哈希会被一并移除,重试会重新上传补齐。 func (s *IngestService) cleanupUploads(ctx context.Context, keys []string) { for _, k := range keys { if d, ok := s.uploader.(storage.Deleter); ok { if err := d.Delete(ctx, k); err != nil { log.Printf("[warn] ingest cleanup: 七牛删除失败 key=%s err=%v", k, err) } } if s.local != nil { if err := s.local.Delete(ctx, k); err != nil { log.Printf("[warn] ingest cleanup: 本地删除失败 key=%s err=%v", k, err) } } } } // fetchedImage 是单张已下载图片的暂存结构,携带去重所需的指纹信息。 type fetchedImage struct { url string key string sha1 string // 内容 sha1(与存储 key 同源) phash string // dHash 的 pgvector 二进制向量串,空串表示无法解码 dupID uint32 // 命中近重复时的参考图 id isDup uint8 // 是否标记为近重复(供审核留痕) } // sqlNull 把 phash 字符串转成可空向量字段:空串 → NULL(不参与近邻检索)。 func sqlNull(ph string) sql.NullString { return sql.NullString{String: ph, Valid: ph != ""} } // dedupImage 判断单张图是否重复: // - 精确重复(content_sha1 已在库)→ 返回 skip=true(不入库该行); // - 近重复(dHash 汉明距离 ≤ 阈值)→ 仍入库,但标记 is_duplicate + dup_of; // - 否则正常入库。 func (s *IngestService) dedupImage(ctx context.Context, kind, sha1hash, phashBits string) (skip bool, dupID uint32, isDup uint8) { // repo 未注入(如离线单测)时跳过去重,不阻断主流程。生产环境 repo 必不为空。 if s.repo == nil { return false, 0, 0 } var tables []string switch kind { case dto.IngestKindRunway: tables = []string{"brand_runway_draft_images", "brand_runway_images"} case dto.IngestKindStreet: tables = []string{"street_snap_draft_images", "street_snap_images"} default: return false, 0, 0 } if exists, err := s.repo.ImageExistsBySha1(ctx, tables, sha1hash); err == nil && exists { return true, 0, 0 } if phashBits != "" { if id, found, err := s.repo.FindNearDuplicateImage(ctx, tables, phashBits, phash.DefaultThreshold); err == nil && found { return false, id, 1 } } return false, 0, 0 } // downloadOne 把单张远程图下载后上传到存储(七牛优先,失败兜底本地), // 返回 (访问地址, 对象 key, 内容 sha1, dHash 向量串, error)。访问地址可直接写入数据库 // (七牛为完整 https URL,本地为相对 /uploads 路径)。 func (s *IngestService) downloadOne(ctx context.Context, u string, idx int, prefix string) (string, string, string, string, error) { resp, err := s.httpClient.Get(u) if err != nil { return "", "", "", "", err } defer resp.Body.Close() if resp.StatusCode != http.StatusOK { return "", "", "", "", fmt.Errorf("status %d", resp.StatusCode) } data, err := io.ReadAll(resp.Body) if err != nil { return "", "", "", "", err } ext := path.Ext(u) if ext == "" || len(ext) > 5 { ext = ".jpg" } // 内容哈希作 key:同一张图(无论来自哪篇文章/job)永远得到相同 key, // 七牛覆盖写即天然幂等,不会因重复采集、崩溃重试、reject 重爬而累积孤儿文件。 h := sha1.Sum(data) hash := hex.EncodeToString(h[:]) key := fmt.Sprintf("%s/%s%s", prefix, hash, ext) ct := resp.Header.Get("Content-Type") // 计算 dHash 指纹(仅用于近重复检索;解码失败返回空串 → NULL,不参与检索)。 ph := phash.Of(data) phashBits := phash.ToVectorBits(ph) // 主上传器(七牛) if url, err := s.uploader.Upload(ctx, key, data, ct); err == nil { return url, key, hash, phashBits, nil } else if s.local != nil { // 兜底本地,避免图片完全丢失 if lurl, lerr := s.local.Upload(ctx, key, data, ct); lerr == nil { return lurl, key, hash, phashBits, nil } else { return "", "", "", "", err } } return "", "", "", "", err }