package ingest import ( "bytes" "context" "crypto/rand" "encoding/hex" "encoding/json" "fmt" "io" "log" "net/http" "strconv" "strings" "time" ) // Client 向后台 ingest 接口发送 HMAC 签名上报的客户端。 type Client struct { Endpoint string Secret string HTTPClient *http.Client } // NewClient 创建上报客户端。endpoint 形如 http://host:8092/admin/internal/ingest。 func NewClient(endpoint, secret string) *Client { return &Client{ Endpoint: endpoint, Secret: secret, HTTPClient: &http.Client{Timeout: 30 * time.Second}, } } // Submit 把一条走秀上报签名后 POST 到后台,返回 job_id(202 入队即返回)。 // 失败返回 error(网络/校验/服务端拒绝),调用方应记日志并继续,不中断整轮抓取。 func (c *Client) Submit(ctx context.Context, p RunwayIngest) (uint32, error) { if c == nil || c.Endpoint == "" { return 0, fmt.Errorf("ingest client 未配置") } body, err := json.Marshal(p) if err != nil { return 0, err } ts := strconv.FormatInt(time.Now().Unix(), 10) nonce, err := newNonce() if err != nil { return 0, err } sig := Sign(c.Secret, ts, nonce, string(body)) req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.Endpoint, bytes.NewReader(body)) if err != nil { return 0, err } req.Header.Set("Content-Type", "application/json") req.Header.Set("X-Signature", sig) req.Header.Set("X-Timestamp", ts) req.Header.Set("X-Nonce", nonce) resp, err := c.HTTPClient.Do(req) if err != nil { return 0, err } defer resp.Body.Close() raw, _ := io.ReadAll(resp.Body) if resp.StatusCode != http.StatusAccepted { return 0, fmt.Errorf("ingest http %d: %s", resp.StatusCode, string(raw)) } var out struct { JobID uint32 `json:"job_id"` Status string `json:"status"` } if err := json.Unmarshal(raw, &out); err != nil { // 202 但响应结构非预期:任务已入队,记日志后当作成功 log.Printf("[Ingest] 响应解析失败但已被接受: %s", string(raw)) return 0, nil } return out.JobID, nil } // newNonce 生成 16 字节随机十六进制串,作为一次性 X-Nonce 防重放。 func newNonce() (string, error) { b := make([]byte, 16) if _, err := rand.Read(b); err != nil { return "", err } return hex.EncodeToString(b), nil } // CrawlBrand 后端抓取任务接口返回的单个品牌(brand_uid 已编码,name 为英文名)。 type CrawlBrand struct { BrandUID string `json:"brand_uid"` Name string `json:"name"` } // crawlBrandsURL 由 ingest endpoint 推导抓取任务接口地址(把末尾 /ingest 换成 /crawl/brands)。 func (c *Client) crawlBrandsURL() string { base := c.Endpoint if i := strings.LastIndex(base, "/ingest"); i != -1 { base = base[:i] } return strings.TrimRight(base, "/") + "/crawl/brands" } // GetCrawlBrands 拉取爬虫要抓取的品牌任务列表(HMAC 签名 GET,复用与上报相同的验签算法)。 // brandID>0 时只取该品牌(对应 --brand 单品牌调试)。返回 brand_uid(直接上送 ingest)+ 英文名。 func (c *Client) GetCrawlBrands(ctx context.Context, brandID uint) ([]CrawlBrand, error) { if c == nil || c.Endpoint == "" { return nil, fmt.Errorf("ingest client 未配置") } url := c.crawlBrandsURL() if brandID > 0 { url += fmt.Sprintf("?brand=%d", brandID) } // GET 无请求体,签名串的 body 部分为空(与后台 IngestAuth 验签一致)。 ts := NowTS() nonce, err := newNonce() if err != nil { return nil, err } sig := Sign(c.Secret, ts, nonce, "") req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil) if err != nil { return nil, err } req.Header.Set("X-Signature", sig) req.Header.Set("X-Timestamp", ts) req.Header.Set("X-Nonce", nonce) resp, err := c.HTTPClient.Do(req) if err != nil { return nil, err } defer resp.Body.Close() raw, _ := io.ReadAll(resp.Body) if resp.StatusCode != http.StatusOK { return nil, fmt.Errorf("crawl brands http %d: %s", resp.StatusCode, string(raw)) } var out struct { Brands []CrawlBrand `json:"brands"` } if err := json.Unmarshal(raw, &out); err != nil { return nil, fmt.Errorf("crawl brands 解析失败: %w", err) } return out.Brands, nil }