Files
spider/internal/ingest/client.go
toom1996 a939848e84 update
2026-09-16 21:31:32 +08:00

172 lines
5.9 KiB
Go
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.

// Package ingest 爬虫 → 后台 ingest 管线的客户端。
//
// 设计(与 backend_v2/internal/middleware/ingest_auth.go + handler/ingest_handler.go 对齐):
// - spider 不再直连数据库:品牌任务经 GET /admin/internal/crawl/brands 拉取,
// 抓取结果经 POST /admin/internal/ingest 上报,由后台 worker 异步下载图、补 season_code、写正式表。
// - 两个接口都走 HMAC-SHA256 验签(X-Signature / X-Timestamp / X-Nonce),
// 签名串 = HMAC_SHA256(secret, timestamp + "." + nonce + "." + bodyRaw)。
// - 时间戳容忍窗口默认 ±300s;nonce 由后台一次性去重防重放。
package ingest
import (
"bytes"
"context"
"crypto/hmac"
"crypto/rand"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"fmt"
"io"
"net/http"
"strconv"
"strings"
"time"
)
// Kind 取值(与 backend_v2/internal/dto/ingest.go 一致)。
const (
KindRunway = "runway" // 走秀(默认,需 brand_uid)
KindStreet = "street" // 街拍(无品牌,需 city/title)
)
// RunwayIngest 上送后端的单场载荷(JSON 字段与 dto.RunwayIngest 完全对齐)。
type RunwayIngest struct {
Kind string `json:"kind"` // runway | street,缺省 runway
BrandUID string `json:"brand_uid"` // 品牌编码 id(hashid),runway 必填
TitleEn string `json:"title_en"` // 英文标题(runway 优先;street 作为单标题)
TitleCn string `json:"title_cn"` // 中文标题(可空)
DescriptionEn string `json:"description_en"` // 英文描述(可空)
DescriptionCn string `json:"description_cn"` // 中文描述(可空)
Year uint16 `json:"year"` // 年份,如 2026
Season string `json:"season"` // spring / fall
CollectionType string `json:"collection_type"` // rtw / menswear / couture / resort / pre_fall
City string `json:"city"` // 地区/城市(street 用)
Images []string `json:"images"` // 铺平的主图 URL 列表(回退用)
Looks []RunwayLook `json:"looks,omitempty"` // 结构化「主图 + 细节图」分组(优先)
}
// RunwayLook 一场秀中的一个 look:1 张主图 + 0..N 张细节图。
type RunwayLook struct {
Main string `json:"main"` // 主图(look)原始 URL
Details []string `json:"details"` // 细节图原始 URL 列表
}
// CrawlBrand 爬虫取任务接口返回的单个品牌。
type CrawlBrand struct {
BrandUID string `json:"brand_uid"`
Name string `json:"name"`
}
// Client 入库管线客户端。
type Client struct {
endpoint string
secret string
http *http.Client
}
// NewClient 创建入库管线客户端。endpoint 形如 http://localhost:8092/admin/internal/ingest。
func NewClient(endpoint, secret string) *Client {
return &Client{
endpoint: endpoint,
secret: secret,
http: &http.Client{Timeout: 30 * time.Second},
}
}
// Submit 把一场走秀/街拍元数据 HMAC 签名后 POST 到后台 ingest,返回 job_id(202 Accepted)。
func (c *Client) Submit(ctx context.Context, p RunwayIngest) (string, error) {
body, err := json.Marshal(p)
if err != nil {
return "", fmt.Errorf("marshal payload: %w", err)
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.endpoint, bytes.NewReader(body))
if err != nil {
return "", fmt.Errorf("new request: %w", err)
}
c.signRequest(req, body)
req.Header.Set("Content-Type", "application/json")
resp, err := c.http.Do(req)
if err != nil {
return "", fmt.Errorf("do request: %w", err)
}
defer resp.Body.Close()
respBody, _ := io.ReadAll(resp.Body)
if resp.StatusCode != http.StatusAccepted {
return "", fmt.Errorf("ingest submit unexpected status %d: %s", resp.StatusCode, string(respBody))
}
var out struct {
JobID string `json:"job_id"`
Status string `json:"status"`
}
_ = json.Unmarshal(respBody, &out)
return out.JobID, nil
}
// GetCrawlBrands 从后台拉取可抓取品牌任务(GET /admin/internal/crawl/brands)。
// brandID>0 时仅返回该品牌(单品牌调试)。
func (c *Client) GetCrawlBrands(ctx context.Context, brandID uint) ([]CrawlBrand, error) {
url := strings.Replace(c.endpoint, "/ingest", "/crawl/brands", 1)
if brandID > 0 {
url += "?brand=" + strconv.FormatUint(uint64(brandID), 10)
}
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
if err != nil {
return nil, fmt.Errorf("new request: %w", err)
}
c.signRequest(req, nil)
resp, err := c.http.Do(req)
if err != nil {
return nil, fmt.Errorf("do request: %w", err)
}
defer resp.Body.Close()
respBody, _ := io.ReadAll(resp.Body)
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("crawl brands unexpected status %d: %s", resp.StatusCode, string(respBody))
}
var out struct {
Brands []CrawlBrand `json:"brands"`
}
if err := json.Unmarshal(respBody, &out); err != nil {
return nil, fmt.Errorf("decode crawl brands: %w", err)
}
return out.Brands, nil
}
// signRequest 对请求体做 HMAC-SHA256 签名并写入验签头。body 为 nil 时按空串处理(GET)。
func (c *Client) signRequest(req *http.Request, body []byte) {
ts := strconv.FormatInt(time.Now().Unix(), 10)
nonce := newNonce()
sig := sign(c.secret, ts, nonce, string(body))
req.Header.Set("X-Signature", sig)
req.Header.Set("X-Timestamp", ts)
req.Header.Set("X-Nonce", nonce)
}
// sign 复刻 backend_v2/internal/pkg/hmac.Sign:HMAC_SHA256(secret, ts + "." + nonce + "." + body)。
func sign(secret, ts, nonce, body string) string {
mac := hmac.New(sha256.New, []byte(secret))
mac.Write([]byte(ts))
mac.Write([]byte("."))
mac.Write([]byte(nonce))
mac.Write([]byte("."))
mac.Write([]byte(body))
return hex.EncodeToString(mac.Sum(nil))
}
// newNonce 生成 16 字节随机十六进制串作为一次性 nonce。
func newNonce() string {
b := make([]byte, 16)
_, _ = rand.Read(b)
return hex.EncodeToString(b)
}