Compare commits

...
7 Commits
Author SHA1 Message Date
19904408334 e47dd99816 feat: 支持流式和非解析HTTP请求及单条DB查询
重构HTTP请求模块,提取doRequestRaw基础方法,新增流式请求(PostStream)、非解析请求(GetNotParse/PostNotParse)及原始响应读取能力;同时在Gfdb接口中新增GetOne单条记录查询方法。
2026-07-23 17:09:27 +08:00
admin 272943d238 运维脚本调整 2026-07-02 11:56:28 +08:00
admin 172c1cc506 common版本更新 2026-06-25 14:55:50 +08:00
admin 915e1f65ec common版本回退 2026-06-24 10:47:47 +08:00
admin c512327358 common版本升级 2026-06-24 09:50:05 +08:00
admin c091ed4984 jaeger修改 2026-06-24 09:41:34 +08:00
admin dc4dd5dc62 common版本升级 2026-06-23 18:17:41 +08:00
5 changed files with 68 additions and 50 deletions
+1
View File
@@ -0,0 +1 @@
.git
+1
View File
@@ -439,6 +439,7 @@ var (
type Gfdb interface { type Gfdb interface {
GetAll(ctx context.Context, sql string, args ...any) (gdb.Result, error) GetAll(ctx context.Context, sql string, args ...any) (gdb.Result, error)
GetOne(ctx context.Context, sql string, args ...any) (gdb.Record, error)
Exec(ctx context.Context, sql string, args ...any) (sql.Result, error) Exec(ctx context.Context, sql string, args ...any) (sql.Result, error)
Model(ctx context.Context, tableNameOrStruct ...any) *model Model(ctx context.Context, tableNameOrStruct ...any) *model
Transaction(ctx context.Context, f func(ctx context.Context, tx gdb.TX) error) error Transaction(ctx context.Context, f func(ctx context.Context, tx gdb.TX) error) error
+1
View File
@@ -1,4 +1,5 @@
module gitea.redpowerfuture.com/red-future/common module gitea.redpowerfuture.com/red-future/common
go 1.26.0 go 1.26.0
require ( require (
+65 -17
View File
@@ -4,12 +4,14 @@ import (
"context" "context"
"errors" "errors"
"fmt" "fmt"
"io"
"net/http" "net/http"
"reflect" "reflect"
"regexp" "regexp"
"strings" "strings"
_ "gitea.redpowerfuture.com/red-future/common/consul" _ "gitea.redpowerfuture.com/red-future/common/consul"
"gitea.redpowerfuture.com/red-future/common/jaeger"
"gitea.redpowerfuture.com/red-future/common/utils" "gitea.redpowerfuture.com/red-future/common/utils"
"github.com/gogf/gf/v2/frame/g" "github.com/gogf/gf/v2/frame/g"
"github.com/gogf/gf/v2/net/gclient" "github.com/gogf/gf/v2/net/gclient"
@@ -68,10 +70,10 @@ func SkipMiddleware(h func(r *ghttp.Request), path string) (handler ghttp.Handle
} }
func RouteRegister(controllers []interface{}) { func RouteRegister(controllers []interface{}) {
//Httpserver.Group("/log", func(group *ghttp.RouterGroup) { Httpserver.Group("/log", func(group *ghttp.RouterGroup) {
// group.Middleware(jaeger.NewTracer) group.Middleware(jaeger.NewTracer)
// group.Bind(controller.OperationLog) //group.Bind(controller.OperationLog)
//}) })
re := regexp.MustCompile("[A-Z]") re := regexp.MustCompile("[A-Z]")
for _, t := range controllers { for _, t := range controllers {
sName := reflect.ValueOf(t).Elem().Type().Name() sName := reflect.ValueOf(t).Elem().Type().Name()
@@ -85,12 +87,9 @@ func RouteRegister(controllers []interface{}) {
go Httpserver.Run() go Httpserver.Run()
} }
// doRequest 统一HTTP请求处理(DELETE用ContentJson发送bodygconv.Struct增加err检查) // doRequestRaw 执行HTTP请求,返回gclient.Response,调用方需自行Close
func doRequest(ctx context.Context, method string, url string, headers map[string]string, target any, data ...any) (err error) { // 统一处理:client克隆、ContentJson设置、请求头注入、GET查询参数转换
err = utils.ValidStructPtr(target) func doRequestRaw(ctx context.Context, method string, url string, headers map[string]string, data ...any) (*gclient.Response, error) {
if err != nil {
return
}
client := Httpclient.Clone() client := Httpclient.Clone()
// POST/PUT/DELETE请求都需要显式用ContentJson序列化body // POST/PUT/DELETE请求都需要显式用ContentJson序列化body
@@ -108,9 +107,9 @@ func doRequest(ctx context.Context, method string, url string, headers map[strin
// 修复:避免data...展开导致的双重包装问题 // 修复:避免data...展开导致的双重包装问题
// 当只有一个元素时,直接传递该元素,避免被包装成数组 // 当只有一个元素时,直接传递该元素,避免被包装成数组
var response *gclient.Response var response *gclient.Response
var err error
// 对于GET请求,将参数转换为map // 对于GET请求,将参数转换为map
if method == http.MethodGet && len(data) > 0 && len(data)%2 == 0 { if method == http.MethodGet && len(data) > 0 && len(data)%2 == 0 {
// 构建query参数map
queryParams := make(map[string]string) queryParams := make(map[string]string)
for i := 0; i < len(data); i += 2 { for i := 0; i < len(data); i += 2 {
if key, ok := data[i].(string); ok && i+1 < len(data) { if key, ok := data[i].(string); ok && i+1 < len(data) {
@@ -124,17 +123,33 @@ func doRequest(ctx context.Context, method string, url string, headers map[strin
} else { } else {
response, err = client.DoRequest(ctx, method, url, data...) response, err = client.DoRequest(ctx, method, url, data...)
} }
return response, err
}
// doRequest 统一HTTP请求处理(同步/异步,解析内部API响应格式并填充target)
func doRequest(ctx context.Context, method string, url string, headers map[string]string, target any, respParse bool, data ...any) (res []byte, err error) {
if target != nil || respParse {
err = utils.ValidStructPtr(target)
if err != nil {
return
}
}
response, err := doRequestRaw(ctx, method, url, headers, data...)
if err != nil { if err != nil {
return return
} }
defer response.Close() defer response.Close()
result := response.ReadAll() result := response.ReadAll()
if !respParse {
return result, nil
}
// 统一处理内部API响应格式:{code:200,message:"",data:{...}} // 统一处理内部API响应格式:{code:200,message:"",data:{...}}
resultStrut := &ghttp.DefaultHandlerResponse{} resultStrut := &ghttp.DefaultHandlerResponse{}
if err = gconv.Struct(result, &resultStrut); err != nil { // 修复:增加err检查 if err = gconv.Struct(result, &resultStrut); err != nil { // 修复:增加err检查
return errors.New("响应解析失败: " + err.Error()) return nil, errors.New("响应解析失败: " + err.Error())
} }
// 添加调试日志:打印解析后的结构 // 添加调试日志:打印解析后的结构
@@ -143,7 +158,7 @@ func doRequest(ctx context.Context, method string, url string, headers map[strin
if resultStrut.Code == 200 || resultStrut.Code == 0 { if resultStrut.Code == 200 || resultStrut.Code == 0 {
if err = gconv.Struct(resultStrut.Data, target); err != nil { // 修复:增加err检查 if err = gconv.Struct(resultStrut.Data, target); err != nil { // 修复:增加err检查
return errors.New("数据解析失败: " + err.Error()) return nil, errors.New("数据解析失败: " + err.Error())
} }
// 添加调试日志:打印最终的target // 添加调试日志:打印最终的target
g.Log().Debugf(ctx, "[HTTP] 最终target: %+v", target) g.Log().Debugf(ctx, "[HTTP] 最终target: %+v", target)
@@ -152,19 +167,52 @@ func doRequest(ctx context.Context, method string, url string, headers map[strin
} }
return return
} }
func Get(ctx context.Context, url string, headers map[string]string, target any, data ...any) (err error) { func Get(ctx context.Context, url string, headers map[string]string, target any, data ...any) (err error) {
err = doRequest(ctx, http.MethodGet, url, headers, target, data...) _, err = doRequest(ctx, http.MethodGet, url, headers, target, true, data...)
return return
} }
func Post(ctx context.Context, url string, headers map[string]string, target any, data ...any) (err error) { func Post(ctx context.Context, url string, headers map[string]string, target any, data ...any) (err error) {
err = doRequest(ctx, http.MethodPost, url, headers, target, data...) _, err = doRequest(ctx, http.MethodPost, url, headers, target, true, data...)
return return
} }
func Put(ctx context.Context, url string, headers map[string]string, target any, data ...any) (err error) { func Put(ctx context.Context, url string, headers map[string]string, target any, data ...any) (err error) {
err = doRequest(ctx, http.MethodPut, url, headers, target, data...) _, err = doRequest(ctx, http.MethodPut, url, headers, target, true, data...)
return return
} }
func Delete(ctx context.Context, url string, headers map[string]string, target any, data ...any) (err error) { func Delete(ctx context.Context, url string, headers map[string]string, target any, data ...any) (err error) {
err = doRequest(ctx, http.MethodDelete, url, headers, target, data...) _, err = doRequest(ctx, http.MethodDelete, url, headers, target, true, data...)
return return
} }
func GetNotParse(ctx context.Context, url string, headers map[string]string, data ...any) (res []byte, err error) {
res, err = doRequest(ctx, http.MethodGet, url, headers, nil, false, data...)
return
}
func PostNotParse(ctx context.Context, url string, headers map[string]string, data ...any) (res []byte, err error) {
res, err = doRequest(ctx, http.MethodPost, url, headers, nil, false, data...)
return
}
// DoStream 流式HTTP请求,返回响应体io.ReadCloser,由调用方自行控制读取和关闭
// 注意:返回的是原始响应流,不会解析内部API响应格式,调用方必须在使用后Close
func doStream(ctx context.Context, method string, url string, headers map[string]string, data ...any) (io.ReadCloser, error) {
response, err := doRequestRaw(ctx, method, url, headers, data...)
if err != nil {
return nil, err
}
// 检查HTTP状态码,提前返回错误
if response.StatusCode < 200 || response.StatusCode >= 300 {
bodyBytes, _ := io.ReadAll(response.Body)
response.Close()
return nil, fmt.Errorf("[HTTP][Stream] 状态码异常: %d, body=%s", response.StatusCode, string(bodyBytes))
}
return response.Body, nil
}
// PostStream POST流式请求,返回响应体io.ReadCloser
func PostStream(ctx context.Context, url string, headers map[string]string, data ...any) (io.ReadCloser, error) {
return doStream(ctx, http.MethodPost, url, headers, data...)
}
-33
View File
@@ -11,12 +11,9 @@ import (
"github.com/gogf/gf/v2/frame/g" "github.com/gogf/gf/v2/frame/g"
"github.com/gogf/gf/v2/net/ghttp" "github.com/gogf/gf/v2/net/ghttp"
"github.com/gogf/gf/v2/net/gtrace" "github.com/gogf/gf/v2/net/gtrace"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/codes" "go.opentelemetry.io/otel/codes"
"go.opentelemetry.io/otel/trace" "go.opentelemetry.io/otel/trace"
"go.opentelemetry.io/otel/trace/embedded"
"go.opentelemetry.io/otel/trace/noop"
) )
var ( var (
@@ -44,10 +41,6 @@ func Init() {
return return
} }
ShutDown = shutdown ShutDown = shutdown
// 包装 TracerProvider:只保留 HTTP Server 追踪,屏蔽 DB/Redis/Client 等内部 span
wrapTracerProvider()
g.Log().Infof(ctx, "✅ Jaeger 初始化成功: %s", jaegerAgent) g.Log().Infof(ctx, "✅ Jaeger 初始化成功: %s", jaegerAgent)
}) })
} }
@@ -57,32 +50,6 @@ func init() {
Init() Init()
} }
// filterTracerProvider 只放行指定 instrument 的 span,其余返回 noop tracer
type filterTracerProvider struct {
embedded.TracerProvider
real trace.TracerProvider
noop trace.TracerProvider
allowed map[string]bool
}
func (f *filterTracerProvider) Tracer(instrumentName string, opts ...trace.TracerOption) trace.Tracer {
if f.allowed[instrumentName] {
return f.real.Tracer(instrumentName, opts...)
}
return f.noop.Tracer(instrumentName, opts...)
}
// wrapTracerProvider 包装全局 TracerProvider,只保留 HTTP Server 追踪
func wrapTracerProvider() {
otel.SetTracerProvider(&filterTracerProvider{
real: otel.GetTracerProvider(),
noop: noop.NewTracerProvider(),
allowed: map[string]bool{
"github.com/gogf/gf/v2/net/ghttp.Server": true,
},
})
}
// NewSpan 创建新的链路追踪 Span // NewSpan 创建新的链路追踪 Span
// spanName: Span 名称,用于在 Jaeger UI 中标识 // spanName: Span 名称,用于在 Jaeger UI 中标识
// 返回带有 Span 的 context 和 Span 对象,调用方需 defer span.End() // 返回带有 Span 的 context 和 Span 对象,调用方需 defer span.End()