Files
bj_power/bj_power_wms/internal/agv/client.go
T

231 lines
6.8 KiB
Go
Raw Normal View History

// Package agv 海康 AGV 调度系统(RCS-2000 V4.2) 对接客户端。
// 能力:下发搬运任务(task/submit)、查询任务状态(task/query)。
// 签名协议详见《AGV海康.md》:HMAC-SHA256 加盐 + MD5 二次摘要,sign 放 URL query。
// 未配置真实 RCS(appKey/appSecret/BaseURL) 或 EnableMock=true 时走 Mock 客户端。
package agv
import (
"bytes"
"context"
"crypto/hmac"
"crypto/md5"
"crypto/rand"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"fmt"
"io"
"net/http"
"sort"
"strings"
"time"
"bj_power_wms/internal/config"
)
// Client AGV 调度客户端抽象
type Client interface {
// SubmitTask 下发搬运任务 fromDock -> toDock,返回海康任务号 robotTaskCode。
SubmitTask(ctx context.Context, fromDock, toDock, carrierCode string) (string, error)
// QueryTask 查询任务状态,返回海康 taskStatus 枚举。
QueryTask(ctx context.Context, robotTaskCode string) (string, error)
}
// New 依据配置返回真实或 Mock 客户端
func New(cfg config.AgvConfig) Client {
if cfg.EnableMock || cfg.BaseURL == "" {
return newMockClient()
}
return newRealClient(cfg)
}
// ---------- 签名与请求(真实海康) ----------
type realClient struct {
cfg config.AgvConfig
hc *http.Client
carrier string
}
func newRealClient(cfg config.AgvConfig) *realClient {
return &realClient{
cfg: cfg,
hc: &http.Client{Timeout: 15 * time.Second},
carrier: cfg.CarrierCode,
}
}
// buildPath 拼接服务前缀 + 业务路径,system
func (c *realClient) url(path string) string {
base := strings.TrimRight(c.cfg.BaseURL, "/")
return base + path
}
// submit 任务下发实现的实际请求
// 海康路径:/rcs/rtas/api/robot/controller/task/submit
func (c *realClient) post(ctx context.Context, path string, payload map[string]any, out any) error {
body, err := json.Marshal(payload)
if err != nil {
return err
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.url(path), bytes.NewReader(body))
if err != nil {
return err
}
c.sign(req, body, path)
resp, err := c.hc.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
raw, err := io.ReadAll(resp.Body)
if err != nil {
return err
}
if resp.StatusCode >= 400 {
return fmt.Errorf("agv http %d: %s", resp.StatusCode, string(raw))
}
var envelope struct {
Code string `json:"code"`
Message string `json:"message"`
Data json.RawMessage `json:"data"`
}
if err := json.Unmarshal(raw, &envelope); err != nil {
return fmt.Errorf("agv parse envelope: %w", err)
}
if envelope.Code != "" && envelope.Code != "SUCCESS" {
return fmt.Errorf("agv code=%s message=%s", envelope.Code, envelope.Message)
}
if out != nil && len(envelope.Data) > 0 {
return json.Unmarshal(envelope.Data, out)
}
return nil
}
// sign 按协议 1.3 节生成鉴权头与 sign 签名。
// 参与签名字段(排序):AUTHORIZATION、HOST、X-LR-APPKEY、X-LR-REQUEST-ID、X-LR-SOURCE、X-LR-TRACE-ID、X-LR-VERSION
// 签名 = MD5( HMAC-SHA256(appSecret, 拼接原串) )sign 作为 URL query 参数。
func (c *realClient) sign(req *http.Request, body []byte, path string) {
nonce := randHex(8)
timestamp := time.Now().Format(time.RFC3339)
requestID := randHex(16)
traceID := randHex(16)
authHeader := fmt.Sprintf("nonce=%q,method=\"HMAC-SHA256\",timestamp=%q", nonce, timestamp)
signKV := map[string]string{
"AUTHORIZATION": authHeader,
"HOST": req.URL.Host,
"X-LR-APPKEY": c.cfg.AppKey,
"X-LR-REQUEST-ID": requestID,
"X-LR-SOURCE": "bj_power_wms",
"X-LR-TRACE-ID": traceID,
"X-LR-VERSION": "v1.0",
}
keys := make([]string, 0, len(signKV))
for k := range signKV {
keys = append(keys, k)
}
sort.Strings(keys)
var sb strings.Builder
for _, k := range keys {
sb.WriteString(k)
sb.WriteString(":")
sb.WriteString(signKV[k])
sb.WriteString("\n")
}
// 请求报文拼接后的原串 = 头拼接串 + 空行 + 消息体
raw := sb.String() + "\n" + string(body)
mac := hmac.New(sha256.New, []byte(c.cfg.AppSecret))
mac.Write([]byte(raw))
h := md5.Sum(mac.Sum(nil))
sign := strings.ToUpper(hex.EncodeToString(h[:]))
q := req.URL.Query()
q.Set("sign", sign)
req.URL.RawQuery = q.Encode()
req.Header.Set("Authorization", authHeader)
req.Header.Set("X-lr-appkey", c.cfg.AppKey)
req.Header.Set("X-lr-request-id", requestID)
req.Header.Set("X-lr-version", "v1.0")
req.Header.Set("X-lr-trace-id", traceID)
req.Header.Set("X-lr-source", "bj_power_wms")
}
// SubmitTask 实现 Client:下发搬运任务。
// 海康路径 /rcs/rtas/api/robot/controller/task/submit
func (c *realClient) SubmitTask(ctx context.Context, fromDock, toDock, carrierCode string) (string, error) {
carrier := carrierCode
if carrier == "" {
carrier = c.carrier
}
payload := map[string]any{
"taskType": "PF-LMR-COMMON",
"targetRoute": []map[string]any{
{"type": "SITE", "code": fromDock, "operation": "COLLECT"},
{"type": "SITE", "code": toDock, "operation": "DELIVERY"},
},
"carrierInfo": []map[string]any{
{"carrierType": "1", "carrierCode": carrier},
},
}
var data struct {
RobotTaskCode string `json:"robotTaskCode"`
Code string `json:"code"`
}
if err := c.post(ctx, "/api/robot/controller/task/submit", payload, &data); err != nil {
return "", err
}
if data.RobotTaskCode == "" {
// 部分版本把任务号放 data.code
if data.Code == "" {
return "", fmt.Errorf("agv submit 未返回任务号")
}
return data.Code, nil
}
return data.RobotTaskCode, nil
}
// QueryTask 实现 Client:查询任务状态。
// 海康路径 /rcs/rtas/api/robot/controller/task/query
func (c *realClient) QueryTask(ctx context.Context, robotTaskCode string) (string, error) {
payload := map[string]any{"robotTaskCode": robotTaskCode}
var data map[string]any
if err := c.post(ctx, "/api/robot/controller/task/query", payload, &data); err != nil {
return "", err
}
st, _ := data["taskStatus"].(string)
if st == "" {
if v, ok := data["status"].(string); ok {
st = v
}
}
if st == "" {
return "QUEUE", nil
}
return st, nil
}
// ---------- Mock 客户端(未接真实 RCS 时用) ----------
type mockClient struct{}
func newMockClient() *mockClient { return &mockClient{} }
func (m *mockClient) SubmitTask(ctx context.Context, fromDock, toDock, carrierCode string) (string, error) {
// 模拟立即成功下发,返回一个可追溯的假任务号(前缀 MOCK-)
return "MOCK-" + randHex(8), nil
}
func (m *mockClient) QueryTask(ctx context.Context, robotTaskCode string) (string, error) {
// 模拟:MOCK 任务一律视为已完成,便于离线联调观察"待搬运→已到位"。
return "FINISHED", nil
}
func randHex(n int) string {
b := make([]byte, n)
if _, err := rand.Read(b); err != nil {
return fmt.Sprintf("%02d", time.Now().UnixNano()%100)
}
return hex.EncodeToString(b)
}