Merge branch 'main' of http://1.117.67.95:3000/team/hyapi-server
This commit is contained in:
@@ -19,8 +19,14 @@ const defaultRequestTimeout = 4 * time.Second
|
||||
|
||||
// queryBillingAPIKeys 查询计费接口:未查得/空结果仍按成功返回空数据,由平台侧计费
|
||||
var queryBillingAPIKeys = map[string]struct{}{
|
||||
"jy000008": {}, // 探针C
|
||||
"jy000017": {}, // 探针A
|
||||
"jy000019": {}, // 商业险有效性
|
||||
"jy000020": {}, // 交强险有效性
|
||||
"jy000022": {}, // 洞侦1.0
|
||||
"jy000028": {}, // 车VIN查车牌号
|
||||
"jy000042": {}, // 借贷意向验证3.0
|
||||
"jy000048": {}, // 申请借贷
|
||||
"jy000052": {}, // 无间司南-纯黑A版
|
||||
}
|
||||
|
||||
|
||||
@@ -5,14 +5,21 @@ import (
|
||||
)
|
||||
|
||||
// generateCurlCommand 生成可直接复现的 curl 命令,便于联调排查。
|
||||
func generateCurlCommand(method, url string, headers map[string]string, body string) string {
|
||||
// proxyURL 非空时附加 -x,便于确认线上是否应走 SOCKS。
|
||||
func generateCurlCommand(method, requestURL string, headers map[string]string, body, proxyURL string) string {
|
||||
var cmd strings.Builder
|
||||
cmd.WriteString("curl -X ")
|
||||
cmd.WriteString(method)
|
||||
cmd.WriteString(" '")
|
||||
cmd.WriteString(url)
|
||||
cmd.WriteString(requestURL)
|
||||
cmd.WriteString("'")
|
||||
|
||||
if p := strings.TrimSpace(proxyURL); p != "" {
|
||||
cmd.WriteString(" \\\n -x '")
|
||||
cmd.WriteString(escapeSingleQuotes(p))
|
||||
cmd.WriteString("'")
|
||||
}
|
||||
|
||||
for key, value := range headers {
|
||||
cmd.WriteString(" \\\n -H '")
|
||||
cmd.WriteString(key)
|
||||
|
||||
47
internal/infrastructure/external/rongxing/http_client_test.go
vendored
Normal file
47
internal/infrastructure/external/rongxing/http_client_test.go
vendored
Normal file
@@ -0,0 +1,47 @@
|
||||
package rongxing
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestNewHTTPClient_Direct(t *testing.T) {
|
||||
client, err := newHTTPClient(5*time.Second, "")
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
if client.Timeout != 5*time.Second {
|
||||
t.Fatalf("timeout = %v", client.Timeout)
|
||||
}
|
||||
if client.Transport != nil && client.Transport != http.DefaultTransport {
|
||||
// 直连允许 Transport 为 nil(使用 DefaultTransport)
|
||||
}
|
||||
}
|
||||
|
||||
func TestNewHTTPClient_Socks5(t *testing.T) {
|
||||
client, err := newHTTPClient(3*time.Second, "socks5://rongxing-vpn:1080")
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
if client.Transport == nil {
|
||||
t.Fatal("expected custom transport for socks5")
|
||||
}
|
||||
}
|
||||
|
||||
func TestNewHTTPClient_HTTPProxy(t *testing.T) {
|
||||
client, err := newHTTPClient(3*time.Second, "http://127.0.0.1:8888")
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
if client.Transport == nil {
|
||||
t.Fatal("expected custom transport for http proxy")
|
||||
}
|
||||
}
|
||||
|
||||
func TestNewHTTPClient_UnsupportedScheme(t *testing.T) {
|
||||
_, err := newHTTPClient(time.Second, "ftp://x")
|
||||
if err == nil {
|
||||
t.Fatal("expected error for unsupported scheme")
|
||||
}
|
||||
}
|
||||
@@ -46,5 +46,6 @@ func NewRongxingServiceWithConfig(cfg *config.Config) (*RongxingService, error)
|
||||
AppID: cfg.Rongxing.AppID,
|
||||
PrivateKey: cfg.Rongxing.PrivateKey,
|
||||
Timeout: timeout,
|
||||
}, logger), nil
|
||||
Proxy: cfg.Rongxing.Proxy,
|
||||
}, logger)
|
||||
}
|
||||
|
||||
@@ -7,7 +7,10 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
@@ -15,6 +18,7 @@ import (
|
||||
"hyapi-server/internal/shared/external_logger"
|
||||
|
||||
"go.uber.org/zap"
|
||||
"golang.org/x/net/proxy"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -32,15 +36,17 @@ type serviceConfig struct {
|
||||
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
|
||||
// loginMu 避免并发登录打爆对方;不缓存 token,每次业务调用重新登录。
|
||||
loginMu sync.Mutex
|
||||
}
|
||||
|
||||
// apiResponse 戎行统一响应。code 可能是数字或字符串;扣费以 consumeFlag 为准。
|
||||
@@ -65,12 +71,52 @@ func (r apiResponse) code() string {
|
||||
}
|
||||
|
||||
// NewRongxingService 创建戎行服务实例
|
||||
func NewRongxingService(cfg serviceConfig, logger *external_logger.ExternalServiceLogger) *RongxingService {
|
||||
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), "/")
|
||||
return &RongxingService{config: cfg, logger: logger}
|
||||
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 获取运行时配置
|
||||
@@ -80,7 +126,7 @@ func (s *RongxingService) GetConfig() serviceConfig {
|
||||
|
||||
// CallAPI 通用业务接口调用。
|
||||
// apiPath 为相对路径(如 /third/loan/info360),reqData 为已组装好的请求体。
|
||||
// Token 获取与 Header 注入由服务内部处理;401/403 时自动刷新 Token 并重试一次。
|
||||
// 每次调用重新登录取 Token(不本地缓存);若业务码/HTTP 仍为 401/403,再登录重试一次。
|
||||
func (s *RongxingService) CallAPI(ctx context.Context, apiPath string, reqData map[string]interface{}) ([]byte, error) {
|
||||
apiKey := strings.Trim(apiPath, "/")
|
||||
|
||||
@@ -118,7 +164,7 @@ func (s *RongxingService) CallAPI(ctx context.Context, apiPath string, reqData m
|
||||
"Content-Type": "application/json",
|
||||
headerDmsToken: token,
|
||||
}
|
||||
curlCmd := generateCurlCommand(http.MethodPost, requestURL, headers, bodyStr)
|
||||
curlCmd := generateCurlCommand(http.MethodPost, requestURL, headers, bodyStr, s.config.Proxy)
|
||||
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost, requestURL, bytes.NewBuffer(bodyBytes))
|
||||
if err != nil {
|
||||
@@ -139,7 +185,6 @@ func (s *RongxingService) CallAPI(ctx context.Context, apiPath string, reqData m
|
||||
respStr := string(respBody)
|
||||
|
||||
if statusCode == http.StatusUnauthorized || statusCode == http.StatusForbidden {
|
||||
s.clearToken()
|
||||
if attempt == 0 {
|
||||
continue
|
||||
}
|
||||
@@ -162,6 +207,11 @@ func (s *RongxingService) CallAPI(ctx context.Context, apiPath string, reqData m
|
||||
}
|
||||
|
||||
code := resp.code()
|
||||
// 对方常以 HTTP 200 + body.code=401 表示 token 无效(非 HTTP 401)
|
||||
if isTokenInvalidCode(code) && attempt == 0 {
|
||||
continue
|
||||
}
|
||||
|
||||
payload := extractBusinessPayload(resp.Data)
|
||||
|
||||
// 扣费只看 consumeFlag:1 扣费(按成功返回),0 不扣费
|
||||
@@ -180,32 +230,20 @@ func (s *RongxingService) CallAPI(ctx context.Context, apiPath string, reqData m
|
||||
return nil, errors.Join(ErrDatasource, errors.New("请求失败"))
|
||||
}
|
||||
|
||||
// getToken 每次重新登录,不缓存 token。
|
||||
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
|
||||
s.loginMu.Lock()
|
||||
defer s.loginMu.Unlock()
|
||||
return s.login(ctx, transactionID)
|
||||
}
|
||||
|
||||
func (s *RongxingService) clearToken() {
|
||||
s.tokenMu.Lock()
|
||||
s.cachedToken = ""
|
||||
s.tokenMu.Unlock()
|
||||
func isTokenInvalidCode(code string) bool {
|
||||
switch strings.TrimSpace(code) {
|
||||
case "401", "403":
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
func (s *RongxingService) login(ctx context.Context, transactionID string) (string, error) {
|
||||
@@ -253,7 +291,7 @@ func (s *RongxingService) login(ctx context.Context, transactionID string) (stri
|
||||
bodyStr := string(bodyBytes)
|
||||
|
||||
headers := map[string]string{"Content-Type": "application/json"}
|
||||
curlCmd := generateCurlCommand(http.MethodPost, requestURL, headers, bodyStr)
|
||||
curlCmd := generateCurlCommand(http.MethodPost, requestURL, headers, bodyStr, s.config.Proxy)
|
||||
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost, requestURL, bytes.NewBuffer(bodyBytes))
|
||||
if err != nil {
|
||||
@@ -309,10 +347,9 @@ func (s *RongxingService) login(ctx context.Context, transactionID string) (stri
|
||||
}
|
||||
|
||||
func (s *RongxingService) doHTTP(req *http.Request) ([]byte, int, error) {
|
||||
client := &http.Client{Timeout: s.config.Timeout}
|
||||
resp, err := client.Do(req)
|
||||
resp, err := s.client.Do(req)
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
return nil, 0, s.wrapTransportError(req.URL.String(), err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
@@ -323,6 +360,49 @@ func (s *RongxingService) doHTTP(req *http.Request) ([]byte, int, error) {
|
||||
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 "网络超时"
|
||||
}
|
||||
// 完整原始错误在同一条日志的 error/cause 字段中,不在别的文件
|
||||
return "其他网络错误(见本条日志 error 全文)"
|
||||
}
|
||||
|
||||
func (s *RongxingService) validateConfig() error {
|
||||
if s.config.BaseURL == "" {
|
||||
return errors.New("戎行 url 未配置")
|
||||
@@ -385,12 +465,35 @@ func (s *RongxingService) logErrorWithCurl(transactionID, apiKey string, err err
|
||||
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)
|
||||
}
|
||||
|
||||
43
internal/infrastructure/external/rongxing/transport_error_test.go
vendored
Normal file
43
internal/infrastructure/external/rongxing/transport_error_test.go
vendored
Normal file
@@ -0,0 +1,43 @@
|
||||
package rongxing
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestClassifyTransportError_Timeout(t *testing.T) {
|
||||
err := errors.New(`Post "http://192.168.3.43:7007/auth/login": context deadline exceeded (Client.Timeout exceeded while awaiting headers)`)
|
||||
got := classifyTransportError(err)
|
||||
if !strings.Contains(got, "请求超时") {
|
||||
t.Fatalf("got %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestClassifyTransportError_Deadline(t *testing.T) {
|
||||
got := classifyTransportError(context.DeadlineExceeded)
|
||||
if !strings.Contains(got, "请求超时") {
|
||||
t.Fatalf("got %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWrapTransportError_IncludesProxy(t *testing.T) {
|
||||
svc, err := NewRongxingService(serviceConfig{
|
||||
BaseURL: "http://192.168.3.43:7007",
|
||||
Timeout: time.Second,
|
||||
Proxy: "socks5://rongxing-vpn:1080",
|
||||
}, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
wrapped := svc.wrapTransportError("http://192.168.3.43:7007/auth/login", context.DeadlineExceeded)
|
||||
msg := wrapped.Error()
|
||||
if !strings.Contains(msg, "proxy=socks5://rongxing-vpn:1080") {
|
||||
t.Fatalf("missing proxy in error: %s", msg)
|
||||
}
|
||||
if !strings.Contains(msg, "diagnosis=") {
|
||||
t.Fatalf("missing diagnosis in error: %s", msg)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user