Files
hyapi-server/internal/infrastructure/external/rongxing/rongxing_service.go
2026-07-23 15:59:00 +08:00

507 lines
14 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 rongxing
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net"
"net/http"
"net/url"
"os"
"strings"
"sync"
"time"
"hyapi-server/internal/shared/external_logger"
"go.uber.org/zap"
"golang.org/x/net/proxy"
)
const (
defaultRequestTimeout = 10 * time.Second
apiKeyLogin = "auth_login"
pathAuthLogin = "/auth/login"
headerDmsToken = "dms-token"
)
// serviceConfig 戎行服务运行时配置
type serviceConfig struct {
BaseURL string
Account string
Password string
AppID string
PrivateKey string
Timeout time.Duration
Proxy string
}
// RongxingService 戎行数据源服务
type RongxingService struct {
config serviceConfig
logger *external_logger.ExternalServiceLogger
client *http.Client
tokenMu sync.RWMutex
cachedToken string
}
// apiResponse 戎行统一响应。code 可能是数字或字符串;扣费以 consumeFlag 为准。
type apiResponse struct {
Flag bool `json:"flag"`
Code json.RawMessage `json:"code"`
Msg string `json:"msg"`
Message string `json:"message"`
Data json.RawMessage `json:"data"`
ConsumeFlag int `json:"consumeFlag"`
}
func (r apiResponse) text() string {
if r.Msg != "" {
return r.Msg
}
return r.Message
}
func (r apiResponse) code() string {
return parseCode(r.Code)
}
// NewRongxingService 创建戎行服务实例
func NewRongxingService(cfg serviceConfig, logger *external_logger.ExternalServiceLogger) (*RongxingService, error) {
if cfg.Timeout <= 0 {
cfg.Timeout = defaultRequestTimeout
}
cfg.BaseURL = strings.TrimRight(strings.TrimSpace(cfg.BaseURL), "/")
client, err := newHTTPClient(cfg.Timeout, cfg.Proxy)
if err != nil {
return nil, err
}
return &RongxingService{config: cfg, logger: logger, client: client}, nil
}
// newHTTPClient 构建 HTTP 客户端proxyURL 支持 socks5/socks5h/http/https空则直连。
func newHTTPClient(timeout time.Duration, proxyURL string) (*http.Client, error) {
client := &http.Client{Timeout: timeout}
proxyURL = strings.TrimSpace(proxyURL)
if proxyURL == "" {
return client, nil
}
u, err := url.Parse(proxyURL)
if err != nil {
return nil, fmt.Errorf("解析 proxy 失败: %w", err)
}
switch strings.ToLower(u.Scheme) {
case "socks5", "socks5h":
dialer, err := proxy.FromURL(u, proxy.Direct)
if err != nil {
return nil, fmt.Errorf("创建 SOCKS 代理失败: %w", err)
}
transport := &http.Transport{}
if cd, ok := dialer.(proxy.ContextDialer); ok {
transport.DialContext = cd.DialContext
} else {
transport.DialContext = func(ctx context.Context, network, addr string) (net.Conn, error) {
return dialer.Dial(network, addr)
}
}
client.Transport = transport
case "http", "https":
client.Transport = &http.Transport{Proxy: http.ProxyURL(u)}
default:
return nil, fmt.Errorf("不支持的 proxy scheme: %s支持 socks5/socks5h/http/https", u.Scheme)
}
return client, nil
}
// GetConfig 获取运行时配置
func (s *RongxingService) GetConfig() serviceConfig {
return s.config
}
// CallAPI 通用业务接口调用。
// apiPath 为相对路径(如 /third/loan/info360reqData 为已组装好的请求体。
// Token 获取与 Header 注入由服务内部处理401/403 时自动刷新 Token 并重试一次。
func (s *RongxingService) CallAPI(ctx context.Context, apiPath string, reqData map[string]interface{}) ([]byte, error) {
apiKey := strings.Trim(apiPath, "/")
var transactionID string
if id, ok := ctx.Value("transaction_id").(string); ok {
transactionID = id
}
if err := s.validateConfig(); err != nil {
err = errors.Join(ErrSystem, err)
s.logError(transactionID, apiKey, "", err, nil)
return nil, err
}
if !strings.HasPrefix(apiPath, "/") {
apiPath = "/" + apiPath
}
requestURL := s.config.BaseURL + apiPath
bodyBytes, err := json.Marshal(reqData)
if err != nil {
err = errors.Join(ErrSystem, err)
s.logError(transactionID, apiKey, "", err, reqData)
return nil, err
}
bodyStr := string(bodyBytes)
for attempt := 0; attempt < 2; attempt++ {
token, err := s.getToken(ctx, transactionID)
if err != nil {
return nil, err
}
headers := map[string]string{
"Content-Type": "application/json",
headerDmsToken: token,
}
curlCmd := generateCurlCommand(http.MethodPost, requestURL, headers, bodyStr, s.config.Proxy)
req, err := http.NewRequestWithContext(ctx, http.MethodPost, requestURL, bytes.NewBuffer(bodyBytes))
if err != nil {
err = errors.Join(ErrSystem, err)
s.logErrorWithCurl(transactionID, apiKey, err, reqData, curlCmd, "")
return nil, err
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set(headerDmsToken, token)
respBody, statusCode, err := s.doHTTP(req)
if err != nil {
err = errors.Join(ErrDatasource, err)
s.logErrorWithCurl(transactionID, apiKey, err, reqData, curlCmd, "")
return nil, err
}
respStr := string(respBody)
if statusCode == http.StatusUnauthorized || statusCode == http.StatusForbidden {
s.clearToken()
if attempt == 0 {
continue
}
err = errors.Join(ErrDatasource, fmt.Errorf("HTTP状态码 %d", statusCode))
s.logErrorWithCurl(transactionID, apiKey, err, reqData, curlCmd, respStr)
return nil, err
}
if statusCode != http.StatusOK {
err = errors.Join(ErrDatasource, fmt.Errorf("HTTP状态码 %d", statusCode))
s.logErrorWithCurl(transactionID, apiKey, err, reqData, curlCmd, respStr)
return nil, err
}
var resp apiResponse
if err := json.Unmarshal(respBody, &resp); err != nil {
err = errors.Join(ErrSystem, fmt.Errorf("响应解析失败: %w, body=%s", err, respStr))
s.logErrorWithCurl(transactionID, apiKey, err, reqData, curlCmd, respStr)
return nil, err
}
code := resp.code()
payload := extractBusinessPayload(resp.Data)
// 扣费只看 consumeFlag1 扣费按成功返回0 不扣费
if IsBillable(resp.ConsumeFlag) {
return payload, nil
}
sentinel := MapNonBillableToErr(code)
err = errors.Join(sentinel, NewRongxingError(code, resp.text()))
if !errors.Is(sentinel, ErrNotFound) {
s.logErrorWithCurl(transactionID, apiKey, err, reqData, curlCmd, respStr)
}
return nil, err
}
return nil, errors.Join(ErrDatasource, errors.New("请求失败"))
}
func (s *RongxingService) getToken(ctx context.Context, transactionID string) (string, error) {
s.tokenMu.RLock()
token := s.cachedToken
s.tokenMu.RUnlock()
if token != "" {
return token, nil
}
s.tokenMu.Lock()
defer s.tokenMu.Unlock()
if s.cachedToken != "" {
return s.cachedToken, nil
}
token, err := s.login(ctx, transactionID)
if err != nil {
return "", err
}
s.cachedToken = token
return token, nil
}
func (s *RongxingService) clearToken() {
s.tokenMu.Lock()
s.cachedToken = ""
s.tokenMu.Unlock()
}
func (s *RongxingService) login(ctx context.Context, transactionID string) (string, error) {
privateKey, err := ParsePrivateKey(s.config.PrivateKey)
if err != nil {
err = errors.Join(ErrSystem, err)
s.logError(transactionID, apiKeyLogin, "", err, nil)
return "", err
}
passwordB64 := EncodePasswordBase64(s.config.Password)
timestamp := time.Now().UnixMilli()
signParams := map[string]interface{}{
"account": s.config.Account,
"password": passwordB64,
"appId": s.config.AppID,
"timestamp": timestamp,
}
content := BuildSignContent(signParams)
sign, err := SignSHA256WithRSA(content, privateKey)
if err != nil {
err = errors.Join(ErrSystem, err)
s.logError(transactionID, apiKeyLogin, "", err, nil)
return "", err
}
payload := map[string]interface{}{
"account": s.config.Account,
"password": passwordB64,
"appId": s.config.AppID,
"timestamp": timestamp,
"sign": sign,
}
requestURL := s.config.BaseURL + pathAuthLogin
bodyBytes, err := json.Marshal(payload)
if err != nil {
err = errors.Join(ErrSystem, err)
s.logError(transactionID, apiKeyLogin, "", err, map[string]interface{}{
"account": s.config.Account,
"appId": s.config.AppID,
})
return "", err
}
bodyStr := string(bodyBytes)
headers := map[string]string{"Content-Type": "application/json"}
curlCmd := generateCurlCommand(http.MethodPost, requestURL, headers, bodyStr, s.config.Proxy)
req, err := http.NewRequestWithContext(ctx, http.MethodPost, requestURL, bytes.NewBuffer(bodyBytes))
if err != nil {
err = errors.Join(ErrSystem, err)
s.logErrorWithCurl(transactionID, apiKeyLogin, err, nil, curlCmd, "")
return "", err
}
req.Header.Set("Content-Type", "application/json")
respBody, statusCode, err := s.doHTTP(req)
if err != nil {
err = errors.Join(ErrDatasource, err)
s.logErrorWithCurl(transactionID, apiKeyLogin, err, nil, curlCmd, "")
return "", err
}
respStr := string(respBody)
if statusCode != http.StatusOK {
err = errors.Join(ErrDatasource, fmt.Errorf("登录 HTTP状态码 %d, body=%s", statusCode, respStr))
s.logErrorWithCurl(transactionID, apiKeyLogin, err, nil, curlCmd, respStr)
return "", err
}
var loginResp apiResponse
if err := json.Unmarshal(respBody, &loginResp); err != nil {
err = errors.Join(ErrSystem, fmt.Errorf("登录响应解析失败: %w, body=%s", err, respStr))
s.logErrorWithCurl(transactionID, apiKeyLogin, err, nil, curlCmd, respStr)
return "", err
}
code := loginResp.code()
if code != CodeSuccess {
rxErr := NewRongxingError(code, loginResp.text())
err = errors.Join(ErrDatasource, rxErr)
s.logErrorWithCurl(transactionID, apiKeyLogin, err, nil, curlCmd, respStr)
return "", err
}
var token string
if err := json.Unmarshal(loginResp.Data, &token); err != nil {
err = errors.Join(ErrSystem, fmt.Errorf("登录响应 Token 解析失败: %w, body=%s", err, respStr))
s.logErrorWithCurl(transactionID, apiKeyLogin, err, nil, curlCmd, respStr)
return "", err
}
token = strings.TrimSpace(token)
if token == "" {
err = errors.Join(ErrSystem, fmt.Errorf("登录响应 Token 为空, body=%s", respStr))
s.logErrorWithCurl(transactionID, apiKeyLogin, err, nil, curlCmd, respStr)
return "", err
}
return token, nil
}
func (s *RongxingService) doHTTP(req *http.Request) ([]byte, int, error) {
resp, err := s.client.Do(req)
if err != nil {
return nil, 0, s.wrapTransportError(req.URL.String(), err)
}
defer resp.Body.Close()
body, err := io.ReadAll(resp.Body)
if err != nil {
return nil, resp.StatusCode, err
}
return body, resp.StatusCode, nil
}
// wrapTransportError 把超时/拒连/代理失败等包装成可读诊断,便于区分「联不通」原因。
func (s *RongxingService) wrapTransportError(targetURL string, err error) error {
diag := classifyTransportError(err)
proxyMode := strings.TrimSpace(s.config.Proxy)
if proxyMode == "" {
proxyMode = "(直连,未配置 proxy)"
}
return fmt.Errorf(
"戎行 HTTP 失败 target=%s proxy=%s diagnosis=%s cause=%w",
targetURL, proxyMode, diag, err,
)
}
func classifyTransportError(err error) string {
if err == nil {
return "unknown"
}
msg := err.Error()
switch {
case os.IsTimeout(err) || errors.Is(err, context.DeadlineExceeded) ||
strings.Contains(msg, "Client.Timeout") || strings.Contains(msg, "deadline exceeded"):
return "请求超时(未在 timeout 内收到响应头;可能:目标 192.168.3.43:7007 不可达、VPN/SOCKS 未转发、或服务无响应)"
case strings.Contains(msg, "connection refused"):
return "连接被拒绝(端口未监听或代理/目标拒绝)"
case strings.Contains(msg, "no such host") || strings.Contains(msg, "lookup"):
return "DNS/主机名解析失败(检查 rongxing-vpn 服务名或目标域名)"
case strings.Contains(msg, "network is unreachable") || strings.Contains(msg, "no route to host"):
return "网络不可达(无路由;直连内网 IP 时常见于未走 VPN/代理)"
case strings.Contains(msg, "i/o timeout") || strings.Contains(msg, "TLS handshake timeout"):
return "传输层超时(链路通但握手/读写超时)"
case strings.Contains(msg, "proxy") || strings.Contains(msg, "socks"):
return "代理链路异常(检查 socks5://rongxing-vpn:1080 与 VPN 容器)"
}
var netErr net.Error
if errors.As(err, &netErr) && netErr.Timeout() {
return "网络超时"
}
return "其他网络错误(详见 cause)"
}
func (s *RongxingService) validateConfig() error {
if s.config.BaseURL == "" {
return errors.New("戎行 url 未配置")
}
if strings.TrimSpace(s.config.Account) == "" {
return errors.New("戎行 account 未配置")
}
if strings.TrimSpace(s.config.Password) == "" {
return errors.New("戎行 password 未配置")
}
if strings.TrimSpace(s.config.AppID) == "" {
return errors.New("戎行 app_id 未配置")
}
if strings.TrimSpace(s.config.PrivateKey) == "" {
return errors.New("戎行 private_key 未配置")
}
return nil
}
func parseCode(raw json.RawMessage) string {
if len(raw) == 0 {
return ""
}
var s string
if err := json.Unmarshal(raw, &s); err == nil {
return strings.TrimSpace(s)
}
var n json.Number
if err := json.Unmarshal(raw, &n); err == nil {
return n.String()
}
return strings.Trim(string(raw), `"`)
}
// extractBusinessPayload 提取对外返回的业务 data。
// 若外层 data 内还嵌套 data如 Info360则取内层标签对象。
func extractBusinessPayload(data json.RawMessage) []byte {
if len(data) == 0 || string(data) == "null" {
return []byte("{}")
}
var wrap struct {
Data json.RawMessage `json:"data"`
}
if err := json.Unmarshal(data, &wrap); err == nil &&
len(wrap.Data) > 0 && string(wrap.Data) != "null" {
return wrap.Data
}
return data
}
func (s *RongxingService) logError(transactionID, apiKey, requestID string, err error, payload interface{}) {
if s.logger == nil {
return
}
s.logger.LogError(requestID, transactionID, apiKey, err, payload)
}
func (s *RongxingService) logErrorWithCurl(transactionID, apiKey string, err error, payload interface{}, curlCmd, respBody string) {
if s.logger == nil {
return
}
proxyMode := strings.TrimSpace(s.config.Proxy)
if proxyMode == "" {
proxyMode = "(直连,未配置 proxy)"
}
s.logger.LogErrorWithFields("rongxing API错误",
zap.String("transaction_id", transactionID),
zap.String("api_code", apiKey),
zap.String("base_url", s.config.BaseURL),
zap.String("proxy", proxyMode),
zap.String("diagnosis", extractDiagnosis(err)),
zap.Error(err),
zap.Any("params", payload),
zap.String("curl", curlCmd),
zap.String("response_body", respBody),
)
}
func extractDiagnosis(err error) string {
if err == nil {
return ""
}
const marker = "diagnosis="
msg := err.Error()
if i := strings.Index(msg, marker); i >= 0 {
rest := msg[i+len(marker):]
if j := strings.Index(rest, " cause="); j >= 0 {
return rest[:j]
}
return rest
}
return classifyTransportError(err)
}