148 lines
4.1 KiB
Go
148 lines
4.1 KiB
Go
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
|
||
}
|