Compare commits

...
27 Commits
Author SHA1 Message Date
19904408334 7a94a208c0 fix(common/http): 修复Authorization头设置时的空指针问题 2026-08-26 14:49:14 +08:00
19904408334 a66a38e074 feat: 新增通用工具框架与WebSocket服务
* 新增 tools 包:统一工具定义、注册表与 Server 接口,对齐 MCP 规范
* 新增 websocket 包:泛化连接管理、心跳、并发写锁与优雅关闭
* 新增参数读取工具函数,避免类型断言静默失败
* 新增 OSS 路径识别与 JSON 扁平映射还原工具
* 修复租户 SQL 条件插入位置,正确处理 GROUP BY 与 ORDER BY 同时出现的场景
2026-08-21 09:26:35 +08:00
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
admin 7e51069595 jaeger修改 2026-06-23 16:51:39 +08:00
19904408334 791c9905df fix: 修正回调地址生成方式 2026-06-22 11:26:49 +08:00
admin 6097209c48 ci/cd调整 2026-06-10 15:01:13 +08:00
admin 1835faddc0 ci/cd调整 2026-06-10 14:54:26 +08:00
19904408334 0ddc2f17b9 fix: 修复snowflake节点并发初始化问题 2026-06-06 14:34:33 +08:00
19904408334 91359f61ac feat: 优化雪花ID生成及添加本地地址工具函数 2026-06-02 19:52:07 +08:00
19904408334 e1829b90bf fix: 修改检查用户限流时候获取用户信息失败状态码为401 2026-04-29 14:36:41 +08:00
admin 6980b31da7 修复consul限流gateway问题 2026-04-29 10:04:30 +08:00
lmk 4960021748 gse配置修改 2026-04-28 14:15:31 +08:00
19904408334 960dca6a50 fix: 修正 creator 字段空值检查逻辑 2026-04-22 15:03:39 +08:00
19904408334 518d45d296 feat: 启用API响应解析逻辑 2026-04-22 14:52:23 +08:00
19904408334 09fe61defd refactor: 移除内部API响应解析逻辑 2026-04-22 13:22:32 +08:00
19904408334 7d2303e5e6 Merge branch 'dev' of http://116.204.74.41:3000/red-future/common into dev 2026-04-22 09:43:59 +08:00
19904408334 38f5d74f4a fix: 修复租户ID和创建者覆盖逻辑 2026-04-22 09:43:18 +08:00
admin c10586ad10 gmq版本 2026-04-21 10:54:10 +08:00
admin 8b7ccb212c gmq版本 2026-04-21 10:29:06 +08:00
admin f5c4977851 gmq版本 2026-04-21 10:27:39 +08:00
admin e051046f77 fix: GSE数据文件从gse/dict目录加载 2026-04-21 10:24:47 +08:00
31 changed files with 1084 additions and 190 deletions
+1
View File
@@ -0,0 +1 @@
.git
+1 -1
View File
@@ -251,7 +251,7 @@ func GetInstanceAddr(ctx context.Context, name string) (addr string, err error)
err = errors.New("获取服务监听器失败") err = errors.New("获取服务监听器失败")
return return
} }
defer watch.Close()
service, err := watch.Proceed() service, err := watch.Proceed()
if err != nil || service == nil { if err != nil || service == nil {
err = errors.New("获取服务实例失败") err = errors.New("获取服务实例失败")
+47 -23
View File
@@ -7,10 +7,11 @@ import (
"fmt" "fmt"
"regexp" "regexp"
"strings" "strings"
"sync"
"time" "time"
"gitea.com/red-future/common/beans" "gitea.redpowerfuture.com/red-future/common/beans"
"gitea.com/red-future/common/utils" "gitea.redpowerfuture.com/red-future/common/utils"
"github.com/bwmarrin/snowflake" "github.com/bwmarrin/snowflake"
"github.com/gogf/gf/v2/crypto/gmd5" "github.com/gogf/gf/v2/crypto/gmd5"
"github.com/gogf/gf/v2/database/gdb" "github.com/gogf/gf/v2/database/gdb"
@@ -27,9 +28,30 @@ import (
// ==================== 缓存管理器(单例) ==================== // ==================== 缓存管理器(单例) ====================
var ( var (
localCache *gcache.Cache localCache *gcache.Cache
snowflakeNode *snowflake.Node
snowflakeOnce sync.Once
) )
func init() {
ctx := context.Background()
snowflakeOnce.Do(func() {
nodeId := genv.Get("APP_NODE", 1).Int64()
// 安全范围 0~1023
if nodeId < 0 || nodeId > 1023 {
nodeId = 1
}
node, err := snowflake.NewNode(nodeId)
if err != nil {
g.Log().Errorf(ctx, "snowflake init failed: %v", err)
return
}
snowflakeNode = node
})
}
// getLocalCache 获取本地缓存实例 // getLocalCache 获取本地缓存实例
func getLocalCache() *gcache.Cache { func getLocalCache() *gcache.Cache {
if localCache == nil { if localCache == nil {
@@ -165,31 +187,30 @@ func insertHook(ctx context.Context, in *gdb.HookInsertInput) (result sql.Result
return nil, err return nil, err
} }
nodeId := genv.Get("APP_NODE", "").Int64() if g.IsEmpty(snowflakeNode) {
if g.IsEmpty(nodeId) { return nil, fmt.Errorf("snowflakeNode is nil")
nodeId = 1
} }
node, err := snowflake.NewNode(nodeId)
if err != nil {
return nil, err
}
for i := range in.Data { for i := range in.Data {
if _, ok := in.Data[i]["id"]; ok { if _, ok := in.Data[i]["id"]; ok {
in.Data[i]["id"] = node.Generate().Int64() in.Data[i]["id"] = snowflakeNode.Generate().Int64()
} }
if _, ok := in.Data[i]["tenant_id"]; ok { if _, ok := in.Data[i]["tenant_id"]; ok {
if !g.IsEmpty(userInfo.TenantId) { if g.IsEmpty(in.Data[i]["tenant_id"]) {
in.Data[i]["tenant_id"] = userInfo.TenantId if !g.IsEmpty(userInfo.TenantId) {
} else { in.Data[i]["tenant_id"] = userInfo.TenantId
return nil, fmt.Errorf("tenantId cannot be empty") } else {
return nil, fmt.Errorf("tenantId cannot be empty")
}
} }
} }
if _, ok := in.Data[i]["creator"]; ok { if _, ok := in.Data[i]["creator"]; ok {
if !g.IsEmpty(userInfo.UserName) { if g.IsEmpty(in.Data[i]["creator"]) {
in.Data[i]["creator"] = userInfo.UserName if !g.IsEmpty(userInfo.UserName) {
} else { in.Data[i]["creator"] = userInfo.UserName
return nil, fmt.Errorf("user info cannot be empty") } else {
return nil, fmt.Errorf("user info cannot be empty")
}
} }
} }
if _, ok := in.Data[i]["updater"]; ok { if _, ok := in.Data[i]["updater"]; ok {
@@ -297,14 +318,16 @@ func selectHook(ctx context.Context, in *gdb.HookSelectInput) (result gdb.Result
return nil, err return nil, err
} }
tenantId = user.TenantId tenantId = user.TenantId
// 【关键修复】找到 SQL 中第一个出现的 ORDER BY / GROUP BY / LIMIT 等关键字位置 // 【关键修复】找到 SQL 中最靠前出现的关键字位置:SQL 子句顺序固定为
// GROUP BY → HAVING → ORDER BY → LIMIT,必须取最早出现的位置,而非关键字列表里先命中的那个。
// 否则同时含 GROUP BY 与 ORDER BY 的查询会命中靠后的 ORDER BY,把 tenant_id 条件拼进 GROUP BY
// (如 GROUP BY DATE(created_at) AND tenant_id = 94),导致 pq: argument of AND must be type boolean。
sql := in.Sql sql := in.Sql
insertPos := len(sql) insertPos := len(sql)
keywords := []string{" ORDER BY ", " GROUP BY ", " HAVING ", " LIMIT ", " FOR UPDATE"} keywords := []string{" GROUP BY ", " HAVING ", " ORDER BY ", " LIMIT ", " FOR UPDATE"}
for _, kw := range keywords { for _, kw := range keywords {
if idx := gstr.PosI(sql, kw); idx != -1 { if idx := gstr.PosI(sql, kw); idx != -1 && idx < insertPos {
insertPos = idx insertPos = idx
break
} }
} }
@@ -418,6 +441,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 -1
View File
@@ -16,7 +16,7 @@ import (
"syscall" "syscall"
"time" "time"
"gitea.com/red-future/common/log/consts" "gitea.redpowerfuture.com/red-future/common/log/consts"
"github.com/gogf/gf/v2/frame/g" "github.com/gogf/gf/v2/frame/g"
"github.com/gogf/gf/v2/os/glog" "github.com/gogf/gf/v2/os/glog"
+4 -4
View File
@@ -11,12 +11,12 @@ import (
"fmt" "fmt"
"time" "time"
"gitea.com/red-future/common/log/consts" "gitea.redpowerfuture.com/red-future/common/log/consts"
"go.mongodb.org/mongo-driver/v2/event" "go.mongodb.org/mongo-driver/v2/event"
"gitea.com/red-future/common/beans" "gitea.redpowerfuture.com/red-future/common/beans"
"gitea.com/red-future/common/log/model/entity" "gitea.redpowerfuture.com/red-future/common/log/model/entity"
"gitea.com/red-future/common/utils" "gitea.redpowerfuture.com/red-future/common/utils"
"github.com/gogf/gf/v2/container/gvar" "github.com/gogf/gf/v2/container/gvar"
"github.com/gogf/gf/v2/errors/gerror" "github.com/gogf/gf/v2/errors/gerror"
"github.com/gogf/gf/v2/frame/g" "github.com/gogf/gf/v2/frame/g"
+2 -2
View File
@@ -10,8 +10,8 @@ import (
"fmt" "fmt"
"time" "time"
"gitea.com/red-future/common/beans" "gitea.redpowerfuture.com/red-future/common/beans"
"gitea.com/red-future/common/utils" "gitea.redpowerfuture.com/red-future/common/utils"
"github.com/gogf/gf/v2/container/gvar" "github.com/gogf/gf/v2/container/gvar"
"github.com/gogf/gf/v2/errors/gerror" "github.com/gogf/gf/v2/errors/gerror"
"github.com/gogf/gf/v2/frame/g" "github.com/gogf/gf/v2/frame/g"
+6 -12
View File
@@ -1,4 +1,4 @@
module gitea.com/red-future/common module gitea.redpowerfuture.com/red-future/common
go 1.26.0 go 1.26.0
@@ -10,15 +10,14 @@ require (
github.com/gogf/gf/contrib/trace/otlphttp/v2 v2.9.5 github.com/gogf/gf/contrib/trace/otlphttp/v2 v2.9.5
github.com/gogf/gf/v2 v2.9.5 github.com/gogf/gf/v2 v2.9.5
github.com/google/uuid v1.6.0 github.com/google/uuid v1.6.0
github.com/gorilla/websocket v1.5.4-0.20250319132907-e064f32e3674
github.com/hashicorp/consul/api v1.26.1 github.com/hashicorp/consul/api v1.26.1
github.com/meilisearch/meilisearch-go v0.36.1 github.com/meilisearch/meilisearch-go v0.36.1
github.com/minio/minio-go/v7 v7.0.97
github.com/nats-io/nats.go v1.48.0
github.com/olivere/elastic/v7 v7.0.32 github.com/olivere/elastic/v7 v7.0.32
github.com/r3labs/diff/v2 v2.15.1 github.com/r3labs/diff/v2 v2.15.1
github.com/rabbitmq/amqp091-go v1.10.0
github.com/rpcxio/rpcx-consul v0.1.1 github.com/rpcxio/rpcx-consul v0.1.1
github.com/smallnest/rpcx v1.9.1 github.com/smallnest/rpcx v1.9.1
github.com/tidwall/sjson v1.2.5
github.com/tiger1103/gfast-token v1.0.10 github.com/tiger1103/gfast-token v1.0.10
go.mongodb.org/mongo-driver/v2 v2.4.0 go.mongodb.org/mongo-driver/v2 v2.4.0
go.opentelemetry.io/otel v1.38.0 go.opentelemetry.io/otel v1.38.0
@@ -55,7 +54,6 @@ require (
github.com/fatih/color v1.18.0 // indirect github.com/fatih/color v1.18.0 // indirect
github.com/fsnotify/fsnotify v1.9.0 // indirect github.com/fsnotify/fsnotify v1.9.0 // indirect
github.com/fxamacker/cbor/v2 v2.9.0 // indirect github.com/fxamacker/cbor/v2 v2.9.0 // indirect
github.com/go-ini/ini v1.67.0 // indirect
github.com/go-logr/logr v1.4.3 // indirect github.com/go-logr/logr v1.4.3 // indirect
github.com/go-logr/stdr v1.2.2 // indirect github.com/go-logr/stdr v1.2.2 // indirect
github.com/go-ole/go-ole v1.2.4 // indirect github.com/go-ole/go-ole v1.2.4 // indirect
@@ -75,7 +73,6 @@ require (
github.com/google/flatbuffers v1.12.1 // indirect github.com/google/flatbuffers v1.12.1 // indirect
github.com/google/gnostic-models v0.7.0 // indirect github.com/google/gnostic-models v0.7.0 // indirect
github.com/google/pprof v0.0.0-20250403155104-27863c87afa6 // indirect github.com/google/pprof v0.0.0-20250403155104-27863c87afa6 // indirect
github.com/gorilla/websocket v1.5.4-0.20250319132907-e064f32e3674 // indirect
github.com/grandcat/zeroconf v1.0.0 // indirect github.com/grandcat/zeroconf v1.0.0 // indirect
github.com/grokify/html-strip-tags-go v0.1.0 // indirect github.com/grokify/html-strip-tags-go v0.1.0 // indirect
github.com/grpc-ecosystem/grpc-gateway/v2 v2.27.2 // indirect github.com/grpc-ecosystem/grpc-gateway/v2 v2.27.2 // indirect
@@ -94,7 +91,6 @@ require (
github.com/kavu/go_reuseport v1.5.0 // indirect github.com/kavu/go_reuseport v1.5.0 // indirect
github.com/klauspost/compress v1.18.0 // indirect github.com/klauspost/compress v1.18.0 // indirect
github.com/klauspost/cpuid/v2 v2.2.11 // indirect github.com/klauspost/cpuid/v2 v2.2.11 // indirect
github.com/klauspost/crc32 v1.3.0 // indirect
github.com/klauspost/reedsolomon v1.12.4 // indirect github.com/klauspost/reedsolomon v1.12.4 // indirect
github.com/libp2p/go-sockaddr v0.2.0 // indirect github.com/libp2p/go-sockaddr v0.2.0 // indirect
github.com/magiconair/properties v1.8.10 // indirect github.com/magiconair/properties v1.8.10 // indirect
@@ -103,15 +99,11 @@ require (
github.com/mattn/go-isatty v0.0.20 // indirect github.com/mattn/go-isatty v0.0.20 // indirect
github.com/mattn/go-runewidth v0.0.16 // indirect github.com/mattn/go-runewidth v0.0.16 // indirect
github.com/miekg/dns v1.1.63 // indirect github.com/miekg/dns v1.1.63 // indirect
github.com/minio/crc64nvme v1.1.0 // indirect
github.com/minio/md5-simd v1.1.2 // indirect
github.com/mitchellh/go-homedir v1.1.0 // indirect github.com/mitchellh/go-homedir v1.1.0 // indirect
github.com/mitchellh/mapstructure v1.5.0 // indirect github.com/mitchellh/mapstructure v1.5.0 // indirect
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect
github.com/modern-go/reflect2 v1.0.3-0.20250322232337-35a7c28c31ee // indirect github.com/modern-go/reflect2 v1.0.3-0.20250322232337-35a7c28c31ee // indirect
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
github.com/nats-io/nkeys v0.4.11 // indirect
github.com/nats-io/nuid v1.0.1 // indirect
github.com/olekukonko/errors v1.1.0 // indirect github.com/olekukonko/errors v1.1.0 // indirect
github.com/olekukonko/ll v0.0.9 // indirect github.com/olekukonko/ll v0.0.9 // indirect
github.com/olekukonko/tablewriter v1.1.0 // indirect github.com/olekukonko/tablewriter v1.1.0 // indirect
@@ -128,7 +120,6 @@ require (
github.com/rivo/uniseg v0.4.7 // indirect github.com/rivo/uniseg v0.4.7 // indirect
github.com/rpcxio/libkv v0.5.1 // indirect github.com/rpcxio/libkv v0.5.1 // indirect
github.com/rs/cors v1.11.1 // indirect github.com/rs/cors v1.11.1 // indirect
github.com/rs/xid v1.6.0 // indirect
github.com/rubyist/circuitbreaker v2.2.1+incompatible // indirect github.com/rubyist/circuitbreaker v2.2.1+incompatible // indirect
github.com/shirou/gopsutil/v3 v3.21.6 // indirect github.com/shirou/gopsutil/v3 v3.21.6 // indirect
github.com/smallnest/quick v0.2.0 // indirect github.com/smallnest/quick v0.2.0 // indirect
@@ -137,6 +128,9 @@ require (
github.com/spf13/pflag v1.0.9 // indirect github.com/spf13/pflag v1.0.9 // indirect
github.com/templexxx/cpufeat v0.0.0-20180724012125-cef66df7f161 // indirect github.com/templexxx/cpufeat v0.0.0-20180724012125-cef66df7f161 // indirect
github.com/templexxx/xor v0.0.0-20191217153810-f85b25db303b // indirect github.com/templexxx/xor v0.0.0-20191217153810-f85b25db303b // indirect
github.com/tidwall/gjson v1.18.0 // indirect
github.com/tidwall/match v1.1.1 // indirect
github.com/tidwall/pretty v1.2.1 // indirect
github.com/tinylib/msgp v1.3.0 // indirect github.com/tinylib/msgp v1.3.0 // indirect
github.com/tjfoc/gmsm v1.4.1 // indirect github.com/tjfoc/gmsm v1.4.1 // indirect
github.com/tklauser/go-sysconf v0.3.6 // indirect github.com/tklauser/go-sysconf v0.3.6 // indirect
+2 -20
View File
@@ -140,8 +140,6 @@ github.com/gkampitakis/go-snaps v0.5.15 h1:amyJrvM1D33cPHwVrjo9jQxX8g/7E2wYdZ+01
github.com/gkampitakis/go-snaps v0.5.15/go.mod h1:HNpx/9GoKisdhw9AFOBT1N7DBs9DiHo/hGheFGBZ+mc= github.com/gkampitakis/go-snaps v0.5.15/go.mod h1:HNpx/9GoKisdhw9AFOBT1N7DBs9DiHo/hGheFGBZ+mc=
github.com/go-ego/gse v1.0.2 h1:+27lYFPhQEhA9igtdOsJPRKYL/k3TwYsxBF5jr6KFv4= github.com/go-ego/gse v1.0.2 h1:+27lYFPhQEhA9igtdOsJPRKYL/k3TwYsxBF5jr6KFv4=
github.com/go-ego/gse v1.0.2/go.mod h1:Fy35G+q7VV7Et1zIKO8o/sW1kkugV3znXap/lF/11zc= github.com/go-ego/gse v1.0.2/go.mod h1:Fy35G+q7VV7Et1zIKO8o/sW1kkugV3znXap/lF/11zc=
github.com/go-ini/ini v1.67.0 h1:z6ZrTEZqSWOTyH2FlglNbNgARyHG8oLW9gMELqKr06A=
github.com/go-ini/ini v1.67.0/go.mod h1:ByCAeIL28uOIIG0E3PJtZPDL8WnHpFKFOtgjp+3Ies8=
github.com/go-kit/kit v0.8.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as= github.com/go-kit/kit v0.8.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as=
github.com/go-kit/kit v0.9.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as= github.com/go-kit/kit v0.9.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as=
github.com/go-kit/kit v0.10.0/go.mod h1:xUsJbQ/Fp4kEt7AFgCuvyX4a71u8h9jB8tj/ORgOZ7o= github.com/go-kit/kit v0.10.0/go.mod h1:xUsJbQ/Fp4kEt7AFgCuvyX4a71u8h9jB8tj/ORgOZ7o=
@@ -353,11 +351,8 @@ github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI
github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck= github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck=
github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo= github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo=
github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ= github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ=
github.com/klauspost/cpuid/v2 v2.0.1/go.mod h1:FInQzS24/EEf25PyTYn52gqo7WaD8xa0213Md/qVLRg=
github.com/klauspost/cpuid/v2 v2.2.11 h1:0OwqZRYI2rFrjS4kvkDnqJkKHdHaRnCm68/DY4OxRzU= github.com/klauspost/cpuid/v2 v2.2.11 h1:0OwqZRYI2rFrjS4kvkDnqJkKHdHaRnCm68/DY4OxRzU=
github.com/klauspost/cpuid/v2 v2.2.11/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= github.com/klauspost/cpuid/v2 v2.2.11/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0=
github.com/klauspost/crc32 v1.3.0 h1:sSmTt3gUt81RP655XGZPElI0PelVTZ6YwCRnPSupoFM=
github.com/klauspost/crc32 v1.3.0/go.mod h1:D7kQaZhnkX/Y0tstFGf8VUzv2UofNGqCjnC3zdHB0Hw=
github.com/klauspost/reedsolomon v1.12.4 h1:5aDr3ZGoJbgu/8+j45KtUJxzYm8k08JGtB9Wx1VQ4OA= github.com/klauspost/reedsolomon v1.12.4 h1:5aDr3ZGoJbgu/8+j45KtUJxzYm8k08JGtB9Wx1VQ4OA=
github.com/klauspost/reedsolomon v1.12.4/go.mod h1:d3CzOMOt0JXGIFZm1StgkyF14EYr3xneR2rNWo7NcMU= github.com/klauspost/reedsolomon v1.12.4/go.mod h1:d3CzOMOt0JXGIFZm1StgkyF14EYr3xneR2rNWo7NcMU=
github.com/konsorten/go-windows-terminal-sequences v1.0.1/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ= github.com/konsorten/go-windows-terminal-sequences v1.0.1/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ=
@@ -415,12 +410,6 @@ github.com/miekg/dns v1.1.27/go.mod h1:KNUDUusw/aVsxyTYZM1oqvCicbwhgbNgztCETuNZ7
github.com/miekg/dns v1.1.41/go.mod h1:p6aan82bvRIyn+zDIv9xYNUpwa73JcSh9BKwknJysuI= github.com/miekg/dns v1.1.41/go.mod h1:p6aan82bvRIyn+zDIv9xYNUpwa73JcSh9BKwknJysuI=
github.com/miekg/dns v1.1.63 h1:8M5aAw6OMZfFXTT7K5V0Eu5YiiL8l7nUAkyN6C9YwaY= github.com/miekg/dns v1.1.63 h1:8M5aAw6OMZfFXTT7K5V0Eu5YiiL8l7nUAkyN6C9YwaY=
github.com/miekg/dns v1.1.63/go.mod h1:6NGHfjhpmr5lt3XPLuyfDJi5AXbNIPM9PY6H6sF1Nfs= github.com/miekg/dns v1.1.63/go.mod h1:6NGHfjhpmr5lt3XPLuyfDJi5AXbNIPM9PY6H6sF1Nfs=
github.com/minio/crc64nvme v1.1.0 h1:e/tAguZ+4cw32D+IO/8GSf5UVr9y+3eJcxZI2WOO/7Q=
github.com/minio/crc64nvme v1.1.0/go.mod h1:eVfm2fAzLlxMdUGc0EEBGSMmPwmXD5XiNRpnu9J3bvg=
github.com/minio/md5-simd v1.1.2 h1:Gdi1DZK69+ZVMoNHRXJyNcxrMA4dSxoYHZSQbirFg34=
github.com/minio/md5-simd v1.1.2/go.mod h1:MzdKDxYpY2BT9XQFocsiZf/NKVtR7nkE4RoEpN+20RM=
github.com/minio/minio-go/v7 v7.0.97 h1:lqhREPyfgHTB/ciX8k2r8k0D93WaFqxbJX36UZq5occ=
github.com/minio/minio-go/v7 v7.0.97/go.mod h1:re5VXuo0pwEtoNLsNuSr0RrLfT/MBtohwdaSmPPSRSk=
github.com/mitchellh/cli v1.0.0/go.mod h1:hNIlj7HEI86fIcpObd7a0FcrxTWetlwJDGcceTlRvqc= github.com/mitchellh/cli v1.0.0/go.mod h1:hNIlj7HEI86fIcpObd7a0FcrxTWetlwJDGcceTlRvqc=
github.com/mitchellh/cli v1.1.0/go.mod h1:xcISNoH86gajksDmfB23e/pu+B+GeFRMYmoHXxx3xhI= github.com/mitchellh/cli v1.1.0/go.mod h1:xcISNoH86gajksDmfB23e/pu+B+GeFRMYmoHXxx3xhI=
github.com/mitchellh/go-homedir v1.0.0/go.mod h1:SfyaCUpYCn1Vlf4IUYiD9fPX4A5wJrkLzIz1N1q0pr0= github.com/mitchellh/go-homedir v1.0.0/go.mod h1:SfyaCUpYCn1Vlf4IUYiD9fPX4A5wJrkLzIz1N1q0pr0=
@@ -450,13 +439,8 @@ github.com/nats-io/jwt v0.3.0/go.mod h1:fRYCDE99xlTsqUzISS1Bi75UBJ6ljOJQOAAu5Vgl
github.com/nats-io/jwt v0.3.2/go.mod h1:/euKqTS1ZD+zzjYrY7pseZrTtWQSjujC7xjPc8wL6eU= github.com/nats-io/jwt v0.3.2/go.mod h1:/euKqTS1ZD+zzjYrY7pseZrTtWQSjujC7xjPc8wL6eU=
github.com/nats-io/nats-server/v2 v2.1.2/go.mod h1:Afk+wRZqkMQs/p45uXdrVLuab3gwv3Z8C4HTBu8GD/k= github.com/nats-io/nats-server/v2 v2.1.2/go.mod h1:Afk+wRZqkMQs/p45uXdrVLuab3gwv3Z8C4HTBu8GD/k=
github.com/nats-io/nats.go v1.9.1/go.mod h1:ZjDU1L/7fJ09jvUSRVBR2e7+RnLiiIQyqyzEE/Zbp4w= github.com/nats-io/nats.go v1.9.1/go.mod h1:ZjDU1L/7fJ09jvUSRVBR2e7+RnLiiIQyqyzEE/Zbp4w=
github.com/nats-io/nats.go v1.48.0 h1:pSFyXApG+yWU/TgbKCjmm5K4wrHu86231/w84qRVR+U=
github.com/nats-io/nats.go v1.48.0/go.mod h1:iRWIPokVIFbVijxuMQq4y9ttaBTMe0SFdlZfMDd+33g=
github.com/nats-io/nkeys v0.1.0/go.mod h1:xpnFELMwJABBLVhffcfd1MZx6VsNRFpEugbxziKVo7w= github.com/nats-io/nkeys v0.1.0/go.mod h1:xpnFELMwJABBLVhffcfd1MZx6VsNRFpEugbxziKVo7w=
github.com/nats-io/nkeys v0.1.3/go.mod h1:xpnFELMwJABBLVhffcfd1MZx6VsNRFpEugbxziKVo7w= github.com/nats-io/nkeys v0.1.3/go.mod h1:xpnFELMwJABBLVhffcfd1MZx6VsNRFpEugbxziKVo7w=
github.com/nats-io/nkeys v0.4.11 h1:q44qGV008kYd9W1b1nEBkNzvnWxtRSQ7A8BoqRrcfa0=
github.com/nats-io/nkeys v0.4.11/go.mod h1:szDimtgmfOi9n25JpfIdGw12tZFYXqhGxjhVxsatHVE=
github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw=
github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c= github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c=
github.com/nxadm/tail v1.4.4/go.mod h1:kenIhsEOeOJmVchQTgglprH7qJGnHDVpk1VPCcaMI8A= github.com/nxadm/tail v1.4.4/go.mod h1:kenIhsEOeOJmVchQTgglprH7qJGnHDVpk1VPCcaMI8A=
github.com/oklog/oklog v0.3.2/go.mod h1:FCV+B7mhrz4o+ueLpx+KqkyXRGMWOYEvfiXtdGtbWGs= github.com/oklog/oklog v0.3.2/go.mod h1:FCV+B7mhrz4o+ueLpx+KqkyXRGMWOYEvfiXtdGtbWGs=
@@ -550,8 +534,6 @@ github.com/quic-go/quic-go v0.49.0 h1:w5iJHXwHxs1QxyBv1EHKuC50GX5to8mJAxvtnttJp9
github.com/quic-go/quic-go v0.49.0/go.mod h1:s2wDnmCdooUQBmQfpUSTCYBl1/D4FcqbULMMkASvR6s= github.com/quic-go/quic-go v0.49.0/go.mod h1:s2wDnmCdooUQBmQfpUSTCYBl1/D4FcqbULMMkASvR6s=
github.com/r3labs/diff/v2 v2.15.1 h1:EOrVqPUzi+njlumoqJwiS/TgGgmZo83619FNDB9xQUg= github.com/r3labs/diff/v2 v2.15.1 h1:EOrVqPUzi+njlumoqJwiS/TgGgmZo83619FNDB9xQUg=
github.com/r3labs/diff/v2 v2.15.1/go.mod h1:I8noH9Fc2fjSaMxqF3G2lhDdC0b+JXCfyx85tWFM9kc= github.com/r3labs/diff/v2 v2.15.1/go.mod h1:I8noH9Fc2fjSaMxqF3G2lhDdC0b+JXCfyx85tWFM9kc=
github.com/rabbitmq/amqp091-go v1.10.0 h1:STpn5XsHlHGcecLmMFCtg7mqq0RnD+zFr4uzukfVhBw=
github.com/rabbitmq/amqp091-go v1.10.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o=
github.com/rcrowley/go-metrics v0.0.0-20181016184325-3113b8401b8a/go.mod h1:bCqnVzQkZxMG4s8nGwiZ5l3QUCyqpo9Y+/ZMZ9VjZe4= github.com/rcrowley/go-metrics v0.0.0-20181016184325-3113b8401b8a/go.mod h1:bCqnVzQkZxMG4s8nGwiZ5l3QUCyqpo9Y+/ZMZ9VjZe4=
github.com/rcrowley/go-metrics v0.0.0-20201227073835-cf1acfcdf475 h1:N/ElC8H3+5XpJzTSTfLsJV/mx9Q9g7kxmchpfZyxgzM= github.com/rcrowley/go-metrics v0.0.0-20201227073835-cf1acfcdf475 h1:N/ElC8H3+5XpJzTSTfLsJV/mx9Q9g7kxmchpfZyxgzM=
github.com/rcrowley/go-metrics v0.0.0-20201227073835-cf1acfcdf475/go.mod h1:bCqnVzQkZxMG4s8nGwiZ5l3QUCyqpo9Y+/ZMZ9VjZe4= github.com/rcrowley/go-metrics v0.0.0-20201227073835-cf1acfcdf475/go.mod h1:bCqnVzQkZxMG4s8nGwiZ5l3QUCyqpo9Y+/ZMZ9VjZe4=
@@ -570,8 +552,6 @@ github.com/rpcxio/rpcx-consul v0.1.1 h1:z/IHpIytgChEuHndWlpo4BY0V0mVBSg/XsKsc54f
github.com/rpcxio/rpcx-consul v0.1.1/go.mod h1:N4SjBS0M9HpVdq3CIIXm/MwWv6euXQ39ybhfH2iTh1c= github.com/rpcxio/rpcx-consul v0.1.1/go.mod h1:N4SjBS0M9HpVdq3CIIXm/MwWv6euXQ39ybhfH2iTh1c=
github.com/rs/cors v1.11.1 h1:eU3gRzXLRK57F5rKMGMZURNdIG4EoAmX8k94r9wXWHA= github.com/rs/cors v1.11.1 h1:eU3gRzXLRK57F5rKMGMZURNdIG4EoAmX8k94r9wXWHA=
github.com/rs/cors v1.11.1/go.mod h1:XyqrcTp5zjWr1wsJ8PIRZssZ8b/WMcMf71DJnit4EMU= github.com/rs/cors v1.11.1/go.mod h1:XyqrcTp5zjWr1wsJ8PIRZssZ8b/WMcMf71DJnit4EMU=
github.com/rs/xid v1.6.0 h1:fV591PaemRlL6JfRxGDEPl69wICngIQ3shQtzfy2gxU=
github.com/rs/xid v1.6.0/go.mod h1:7XoLgs4eV+QndskICGsho+ADou8ySMSjJKDIan90Nz0=
github.com/rubyist/circuitbreaker v2.2.1+incompatible h1:KUKd/pV8Geg77+8LNDwdow6rVCAYOp8+kHUyFvL6Mhk= github.com/rubyist/circuitbreaker v2.2.1+incompatible h1:KUKd/pV8Geg77+8LNDwdow6rVCAYOp8+kHUyFvL6Mhk=
github.com/rubyist/circuitbreaker v2.2.1+incompatible/go.mod h1:Ycs3JgJADPuzJDwffe12k6BZT8hxVi6lFK+gWYJLN4A= github.com/rubyist/circuitbreaker v2.2.1+incompatible/go.mod h1:Ycs3JgJADPuzJDwffe12k6BZT8hxVi6lFK+gWYJLN4A=
github.com/russross/blackfriday/v2 v2.0.1/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM= github.com/russross/blackfriday/v2 v2.0.1/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM=
@@ -628,10 +608,12 @@ github.com/templexxx/cpufeat v0.0.0-20180724012125-cef66df7f161 h1:89CEmDvlq/F7S
github.com/templexxx/cpufeat v0.0.0-20180724012125-cef66df7f161/go.mod h1:wM7WEvslTq+iOEAMDLSzhVuOt5BRZ05WirO+b09GHQU= github.com/templexxx/cpufeat v0.0.0-20180724012125-cef66df7f161/go.mod h1:wM7WEvslTq+iOEAMDLSzhVuOt5BRZ05WirO+b09GHQU=
github.com/templexxx/xor v0.0.0-20191217153810-f85b25db303b h1:fj5tQ8acgNUr6O8LEplsxDhUIe2573iLkJc+PqnzZTI= github.com/templexxx/xor v0.0.0-20191217153810-f85b25db303b h1:fj5tQ8acgNUr6O8LEplsxDhUIe2573iLkJc+PqnzZTI=
github.com/templexxx/xor v0.0.0-20191217153810-f85b25db303b/go.mod h1:5XA7W9S6mni3h5uvOC75dA3m9CCCaS83lltmc0ukdi4= github.com/templexxx/xor v0.0.0-20191217153810-f85b25db303b/go.mod h1:5XA7W9S6mni3h5uvOC75dA3m9CCCaS83lltmc0ukdi4=
github.com/tidwall/gjson v1.14.2/go.mod h1:/wbyibRr2FHMks5tjHJ5F8dMZh3AcwJEMf5vlfC0lxk=
github.com/tidwall/gjson v1.18.0 h1:FIDeeyB800efLX89e5a8Y0BNH+LOngJyGrIWxG2FKQY= github.com/tidwall/gjson v1.18.0 h1:FIDeeyB800efLX89e5a8Y0BNH+LOngJyGrIWxG2FKQY=
github.com/tidwall/gjson v1.18.0/go.mod h1:/wbyibRr2FHMks5tjHJ5F8dMZh3AcwJEMf5vlfC0lxk= github.com/tidwall/gjson v1.18.0/go.mod h1:/wbyibRr2FHMks5tjHJ5F8dMZh3AcwJEMf5vlfC0lxk=
github.com/tidwall/match v1.1.1 h1:+Ho715JplO36QYgwN9PGYNhgZvoUSc9X2c80KVTi+GA= github.com/tidwall/match v1.1.1 h1:+Ho715JplO36QYgwN9PGYNhgZvoUSc9X2c80KVTi+GA=
github.com/tidwall/match v1.1.1/go.mod h1:eRSPERbgtNPcGhD8UCthc6PmLEQXEWd3PRB5JTxsfmM= github.com/tidwall/match v1.1.1/go.mod h1:eRSPERbgtNPcGhD8UCthc6PmLEQXEWd3PRB5JTxsfmM=
github.com/tidwall/pretty v1.2.0/go.mod h1:ITEVvHYasfjBbM0u2Pg8T2nJnzm8xPwvNhhsoaGGjNU=
github.com/tidwall/pretty v1.2.1 h1:qjsOFOWWQl+N3RsoF5/ssm1pHmJJwhjlSbZ51I6wMl4= github.com/tidwall/pretty v1.2.1 h1:qjsOFOWWQl+N3RsoF5/ssm1pHmJJwhjlSbZ51I6wMl4=
github.com/tidwall/pretty v1.2.1/go.mod h1:ITEVvHYasfjBbM0u2Pg8T2nJnzm8xPwvNhhsoaGGjNU= github.com/tidwall/pretty v1.2.1/go.mod h1:ITEVvHYasfjBbM0u2Pg8T2nJnzm8xPwvNhhsoaGGjNU=
github.com/tidwall/sjson v1.2.5 h1:kLy8mja+1c9jlljvWTlSazM7cKDRfJuR/bOJhcY5NcY= github.com/tidwall/sjson v1.2.5 h1:kLy8mja+1c9jlljvWTlSazM7cKDRfJuR/bOJhcY5NcY=
+69 -23
View File
@@ -4,14 +4,15 @@ import (
"context" "context"
"errors" "errors"
"fmt" "fmt"
"io"
"net/http" "net/http"
"reflect" "reflect"
"regexp" "regexp"
"strings" "strings"
_ "gitea.com/red-future/common/consul" _ "gitea.redpowerfuture.com/red-future/common/consul"
"gitea.com/red-future/common/jaeger" "gitea.redpowerfuture.com/red-future/common/jaeger"
"gitea.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"
"github.com/gogf/gf/v2/net/ghttp" "github.com/gogf/gf/v2/net/ghttp"
@@ -69,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()
@@ -80,19 +81,15 @@ func RouteRegister(controllers []interface{}) {
return fmt.Sprintf("/%s", strings.ToLower(s)) return fmt.Sprintf("/%s", strings.ToLower(s))
}) })
Httpserver.Group(convertedStr, func(group *ghttp.RouterGroup) { Httpserver.Group(convertedStr, func(group *ghttp.RouterGroup) {
group.Middleware(jaeger.NewTracer)
group.Bind(t) group.Bind(t)
}) })
} }
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
@@ -103,16 +100,16 @@ func doRequest(ctx context.Context, method string, url string, headers map[strin
// 最后设置headers,确保不会被ContentJson覆盖 // 最后设置headers,确保不会被ContentJson覆盖
if len(headers) > 0 { if len(headers) > 0 {
client.SetHeaderMap(headers) client.SetHeaderMap(headers)
} else { } else if r := g.RequestFromCtx(ctx); r != nil {
client.SetHeader("Authorization", g.RequestFromCtx(ctx).GetHeader("Authorization")) client.SetHeader("Authorization", r.GetHeader("Authorization"))
} }
// 修复:避免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) {
@@ -126,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())
} }
// 添加调试日志:打印解析后的结构 // 添加调试日志:打印解析后的结构
@@ -145,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)
@@ -154,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...)
}
+2 -2
View File
@@ -3,8 +3,8 @@ package controller
import ( import (
"context" "context"
"gitea.com/red-future/common/log/model/dto" "gitea.redpowerfuture.com/red-future/common/log/model/dto"
"gitea.com/red-future/common/log/service" "gitea.redpowerfuture.com/red-future/common/log/service"
) )
type operationLog struct{} type operationLog struct{}
+5 -5
View File
@@ -3,15 +3,15 @@ package dao
import ( import (
"context" "context"
"gitea.com/red-future/common/beans" "gitea.redpowerfuture.com/red-future/common/beans"
"gitea.com/red-future/common/db/mongo" "gitea.redpowerfuture.com/red-future/common/db/mongo"
"strings" "strings"
"time" "time"
"gitea.com/red-future/common/log/consts" "gitea.redpowerfuture.com/red-future/common/log/consts"
"gitea.com/red-future/common/log/model/dto" "gitea.redpowerfuture.com/red-future/common/log/model/dto"
"gitea.com/red-future/common/log/model/entity" "gitea.redpowerfuture.com/red-future/common/log/model/entity"
"go.mongodb.org/mongo-driver/v2/bson" "go.mongodb.org/mongo-driver/v2/bson"
) )
+1 -1
View File
@@ -1,7 +1,7 @@
package dto package dto
import ( import (
"gitea.com/red-future/common/beans" "gitea.redpowerfuture.com/red-future/common/beans"
"github.com/gogf/gf/v2/frame/g" "github.com/gogf/gf/v2/frame/g"
"github.com/gogf/gf/v2/os/gtime" "github.com/gogf/gf/v2/os/gtime"
) )
+1 -1
View File
@@ -1,7 +1,7 @@
package entity package entity
import ( import (
"gitea.com/red-future/common/beans" "gitea.redpowerfuture.com/red-future/common/beans"
) )
// OperationLog 操作日志实体 - 用于记录数据增删改操作行为 // OperationLog 操作日志实体 - 用于记录数据增删改操作行为
+5 -5
View File
@@ -2,11 +2,11 @@ package service
import ( import (
"context" "context"
"gitea.com/red-future/common/beans" "gitea.redpowerfuture.com/red-future/common/beans"
"gitea.com/red-future/common/log/dao" "gitea.redpowerfuture.com/red-future/common/log/dao"
"gitea.com/red-future/common/log/model/dto" "gitea.redpowerfuture.com/red-future/common/log/model/dto"
logEntity "gitea.com/red-future/common/log/model/entity" logEntity "gitea.redpowerfuture.com/red-future/common/log/model/entity"
"gitea.com/red-future/common/utils" "gitea.redpowerfuture.com/red-future/common/utils"
"github.com/gogf/gf/v2/util/gconv" "github.com/gogf/gf/v2/util/gconv"
) )
+1 -1
View File
@@ -9,7 +9,7 @@ import (
"sync/atomic" "sync/atomic"
"time" "time"
"gitea.com/red-future/common/utils" "gitea.redpowerfuture.com/red-future/common/utils"
"github.com/alibaba/sentinel-golang/api" "github.com/alibaba/sentinel-golang/api"
"github.com/alibaba/sentinel-golang/core/circuitbreaker" "github.com/alibaba/sentinel-golang/core/circuitbreaker"
"github.com/gogf/gf/v2/frame/g" "github.com/gogf/gf/v2/frame/g"
+3 -3
View File
@@ -4,9 +4,9 @@ import (
"context" "context"
"encoding/json" "encoding/json"
"fmt" "fmt"
"gitea.com/red-future/common/beans" "gitea.redpowerfuture.com/red-future/common/beans"
commonHttp "gitea.com/red-future/common/http" commonHttp "gitea.redpowerfuture.com/red-future/common/http"
"gitea.com/red-future/common/utils" "gitea.redpowerfuture.com/red-future/common/utils"
"github.com/gogf/gf/v2/database/gredis" "github.com/gogf/gf/v2/database/gredis"
"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"
+2 -2
View File
@@ -5,7 +5,7 @@ import (
"fmt" "fmt"
"strings" "strings"
"gitea.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/ghttp" "github.com/gogf/gf/v2/net/ghttp"
"github.com/gogf/gf/v2/util/gconv" "github.com/gogf/gf/v2/util/gconv"
@@ -92,7 +92,7 @@ func UserLimiter(r *ghttp.Request) {
var userName string var userName string
user, err := utils.GetUserInfo(r.GetCtx()) user, err := utils.GetUserInfo(r.GetCtx())
if err != nil { if err != nil {
r.Response.WriteStatusExit(429, err.Error()) r.Response.WriteStatusExit(401, err.Error())
return return
} }
userName = gconv.String(user.UserName) userName = gconv.String(user.UserName)
+2 -2
View File
@@ -2,8 +2,8 @@ package swagger
import ( import (
"fmt" "fmt"
"gitea.com/red-future/common/consul" "gitea.redpowerfuture.com/red-future/common/consul"
"gitea.com/red-future/common/http" "gitea.redpowerfuture.com/red-future/common/http"
"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/util/gconv" "github.com/gogf/gf/v2/util/gconv"
+47
View File
@@ -0,0 +1,47 @@
package tools
import "github.com/gogf/gf/v2/util/gconv"
// 标准入参读取。LLM/外部传入的 map[string]any 中数值可能是 float64、int 或字符串,
// 直接类型断言极易静默失败,统一走 gconv 强转。
// HasArg 判断参数是否存在且非 nil
func HasArg(args map[string]any, key string) bool {
v, ok := args[key]
return ok && v != nil
}
// ArgString 读取字符串参数
func ArgString(args map[string]any, key string) string {
return gconv.String(args[key])
}
// ArgInt 读取整数参数(自动兼容 float64/int/字符串)
func ArgInt(args map[string]any, key string) int {
return gconv.Int(args[key])
}
// ArgFloat64 读取浮点参数
func ArgFloat64(args map[string]any, key string) float64 {
return gconv.Float64(args[key])
}
// ArgBool 读取布尔参数
func ArgBool(args map[string]any, key string) bool {
return gconv.Bool(args[key])
}
// ArgSlice 读取切片参数
func ArgSlice(args map[string]any, key string) []any {
return gconv.SliceAny(args[key])
}
// ArgStrings 读取字符串切片参数
func ArgStrings(args map[string]any, key string) []string {
return gconv.Strings(args[key])
}
// ArgMap 读取对象参数
func ArgMap(args map[string]any, key string) map[string]any {
return gconv.Map(args[key])
}
+20
View File
@@ -0,0 +1,20 @@
// Package tools 提供共享的模型工具框架:工具定义、注册表与结构化返回。
//
// 数据模型对齐 MCP 的 Tool 定义(name + description + inputSchema)。
// 工具是**给模型 function calling 用**的通用能力:模型在推理中动态决定是否调用,
// 消费方依赖 Server 接口(List/Call)而非注册表本身,未来可无缝替换为远程 MCP Server 客户端。
//
// # 分层
//
// - tools(本包):Tool 定义 + 全局注册表 Register + Server 接口,纯框架,不依赖任何业务包
// - 业务侧:各服务在自身代码中定义**通用**工具(含执行实现 Func),通过 init() 注册进共享注册表;
// 工具的"使用"(如 ReAct 执行循环)也由业务服务基于 Server.Call 自行编排
//
// 非通用、绑定到固定业务场景的处理逻辑(如工作流节点的前后置钩子)**不属于工具**,
// 应由各业务在自己的编排层维护,不进入本注册表。
//
// # 消费路径
//
// 注册表内工具由各业务服务的模型调用方消费:模型 function calling 模式下
// 经 Server.List 获取全部工具定义交给模型,按返回的 ToolCall 经 Server.Call 执行。
package tools
+29
View File
@@ -0,0 +1,29 @@
package tools
import "fmt"
// ToolResult 工具统一返回结构。Code=0 表示成功,非 0 为业务/内部错误码;
// 调用方(工作流前置/后置钩子、模型 function calling)统一按 Code 判断结果。
type ToolResult struct {
Code int `json:"code"` // 0=成功,非 0=失败
Message string `json:"message"` // 成功说明或错误信息
Data any `json:"data,omitempty"`
}
// 标准错误码。业务工具可自定义扩展 >500 的错误码。
const (
CodeOK = 0 // 成功
CodeInvalidArgs = 400 // 入参缺失或格式错误
CodeNotFound = 404 // 工具/数据不存在
CodeInternal = 500 // 内部执行错误
)
// OK 构造成功结果
func OK(data any) ToolResult {
return ToolResult{Code: CodeOK, Data: data}
}
// Fail 构造失败结果,message 支持 fmt 格式化
func Fail(code int, format string, args ...any) ToolResult {
return ToolResult{Code: code, Message: fmt.Sprintf(format, args...)}
}
+47
View File
@@ -0,0 +1,47 @@
package tools
import (
"context"
"sort"
)
// Server 提供工具的发现与执行能力,对应 MCP 中 Server 侧的 tools/list 与 tools/call。
// 消费方(工作流前置/后置钩子、/tool/list 接口、模型 function calling)依赖该接口而非注册表本身,
// 未来可无缝替换为远程 MCP Server 客户端实现。
type Server interface {
// List 返回全部已注册工具定义,按名称排序
List(ctx context.Context) ([]*Tool, error)
// Call 按名称调用工具。业务失败通过返回结果的 Code 表达,
// error 仅用于基础设施异常(如 ctx 取消)。
Call(ctx context.Context, name string, args map[string]any) (ToolResult, error)
}
// localServer 基于本地注册表的 Server 实现
type localServer struct{}
// Default 默认本地 Server。未来接远程 MCP 时,替换该变量即可,业务侧零改动。
var Default Server = localServer{}
func (localServer) List(ctx context.Context) ([]*Tool, error) {
list := make([]*Tool, 0, len(registry))
for _, t := range registry {
list = append(list, t)
}
sort.Slice(list, func(i, j int) bool { return list[i].Name < list[j].Name })
return list, nil
}
func (localServer) Call(ctx context.Context, name string, args map[string]any) (ToolResult, error) {
if err := ctx.Err(); err != nil {
return ToolResult{}, err
}
tool := registry[name]
if tool == nil || tool.Func == nil {
return Fail(CodeNotFound, "工具[%s]不存在或未实现", name), nil
}
res, err := tool.Func(ctx, args)
if err != nil {
return Fail(CodeInternal, "工具[%s]执行失败: %v", name, err), nil
}
return res, nil
}
+26
View File
@@ -0,0 +1,26 @@
package tools
import "context"
// Tool 对应 MCP 的 Tool 数据结构:name + description + inputSchema。
// 模型 function calling 通过 InputSchema 感知入参,推理中按需调用 Func。
// 若未来接远程 MCP Server,Func 由客户端统一实现,本结构不变。
type Tool struct {
Name string // 工具唯一标识
Description string // 用途说明(模型必读,用于决定何时调用)
Parameters map[string]any // 入参 JSON Schema
Func func(ctx context.Context, args map[string]any) (ToolResult, error)
}
// registry 工具注册表
var registry = make(map[string]*Tool)
// Register 注册工具,同名覆盖
func Register(list ...*Tool) {
for _, t := range list {
if t == nil || t.Name == "" {
continue
}
registry[t.Name] = t
}
}
+59 -60
View File
@@ -42,78 +42,77 @@ type gseTool struct {
func newGseTool() (tool *gseTool, err error) { func newGseTool() (tool *gseTool, err error) {
// 1. 初始化分词器 // 1. 初始化分词器
var seg gse.Segmenter var seg gse.Segmenter
// 内置词典(无外部文件)
err = seg.LoadDictEmbed()
if err != nil {
return
}
// 内置停用词(v1.0.2 标准)
err = seg.LoadStopEmbed()
if err != nil {
return
}
// 获取GSE数据文件路径 // 2. 初始化 TF-IDF 提取器
gseDataPath := os.Getenv("GSE_DATA_PATH") tfidf := &extracker.TagExtracter{}
tfidf.WithGse(seg)
if gseDataPath != "" { // 尝试从默认路径加载 IDF 字典
// 使用外部数据文件 idfPath := getIdfDictPath()
dictPath := filepath.Join(gseDataPath, "dict", "zh") if idfPath != "" {
idfPath := filepath.Join(gseDataPath, "dict", "zh", "idf.txt") // 如果找到自定义路径,使用 LoadDict 方法加载
stopPath := filepath.Join(gseDataPath, "dict", "zh", "stop.txt") err = tfidf.LoadDict(idfPath)
// 加载词典
err = seg.LoadDict(filepath.Join(dictPath, "dict.txt"))
if err != nil { if err != nil {
return glog.Warningf(context.Background(), "加载自定义 IDF 字典失败 [%s]: %v,将使用默认字典", idfPath, err)
} // 回退到默认加载方式
err = tfidf.LoadIdf()
// 加载停用词 } else {
err = seg.LoadStop(stopPath) glog.Infof(context.Background(), "成功加载自定义 IDF 字典: %s", idfPath)
if err != nil {
glog.Warning(context.Background(), "加载停用词失败,继续:", err)
}
// 2. 初始化 TF-IDF 提取器
tfidf := &extracker.TagExtracter{}
tfidf.WithGse(seg)
err = tfidf.LoadIdf(idfPath)
if err != nil {
return
}
// 3. 初始化 TextRank 提取器
tr := &extracker.TextRanker{}
tr.WithGse(seg)
tool = &gseTool{
seg: seg,
tfidf: tfidf,
tr: tr,
} }
} else { } else {
// 使用内置embed数据 // 使用默认的 IDF 字典
err = seg.LoadDictEmbed()
if err != nil {
return
}
// 内置停用词(v1.0.2 标准)
err = seg.LoadStopEmbed()
if err != nil {
return
}
// 2. 初始化 TF-IDF 提取器
tfidf := &extracker.TagExtracter{}
tfidf.WithGse(seg)
err = tfidf.LoadIdf() err = tfidf.LoadIdf()
if err != nil { }
return
}
// 3. 初始化 TextRank 提取器 if err != nil {
tr := &extracker.TextRanker{} return
tr.WithGse(seg) }
tool = &gseTool{ // 3. 初始化 TextRank 提取器
seg: seg, tr := &extracker.TextRanker{}
tfidf: tfidf, tr.WithGse(seg)
tr: tr,
} tool = &gseTool{
seg: seg,
tfidf: tfidf,
tr: tr,
} }
return return
} }
// getIdfDictPath 获取 IDF 字典文件路径
func getIdfDictPath() string {
// 1. 尝试从容器内的默认挂载路径加载(Docker 卷映射)
containerPath := "/app/dict/zh/idf.txt"
if _, err := os.Stat(containerPath); err == nil {
return containerPath
}
// 2. 尝试从当前工作目录的 dict/zh/idf.txt 加载
workDir, err := os.Getwd()
if err != nil {
return ""
}
localPath := filepath.Join(workDir, "dict", "zh", "idf.txt")
if _, err := os.Stat(localPath); err == nil {
return localPath
}
// 3. 如果没有找到自定义路径,返回空字符串,使用默认字典
return ""
}
// Cut 分词(关键词提取唯一正确模式:精确模式 + HMM) // Cut 分词(关键词提取唯一正确模式:精确模式 + HMM)
func (k *gseTool) Cut(text string) []string { func (k *gseTool) Cut(text string) []string {
return k.seg.Cut(text, true) return k.seg.Cut(text, true)
+43
View File
@@ -0,0 +1,43 @@
package utils
import (
"encoding/json"
"fmt"
"github.com/tidwall/sjson"
)
// IsFlatMap 递归判断 map 是否扁平化
func IsFlatMap(m map[string]interface{}) bool {
for _, v := range m {
switch val := v.(type) {
case map[string]interface{}:
return false
case []interface{}:
for _, item := range val {
if _, ok := item.(map[string]interface{}); ok {
return false
}
}
}
}
return true
}
// UnFlatBySjson 将扁平路径映射还原为嵌套 JSON
func UnFlatBySjson(flatMap map[string]interface{}) (map[string]interface{}, error) {
raw := "{}"
for path, val := range flatMap {
var err error
raw, err = sjson.Set(raw, path, val)
if err != nil {
return nil, fmt.Errorf("sjson set path %s failed: %w", path, err)
}
}
var result map[string]interface{}
if err := json.Unmarshal([]byte(raw), &result); err != nil {
return nil, fmt.Errorf("parse final json failed: %w", err)
}
return result, nil
}
+41
View File
@@ -0,0 +1,41 @@
package utils
import (
"context"
"fmt"
"regexp"
"github.com/gogf/gf/v2/frame/g"
)
// ossObjectPathPattern 匹配 MinIO 上传生成的对象路径(不带 http 前缀的相对路径):
// /YYYY-MM-DD/32位uuid.扩展名,如 /2026-08-19/1e9d9e48-3f6b-4a2c-8d5e-1f2a3b4c.png
var ossObjectPathPattern = regexp.MustCompile(`^/\d{4}-\d{2}-\d{2}/[0-9a-fA-F-]{32}\.[a-zA-Z0-9]{1,10}$`)
// IsOSSPath 判断字符串是否为 MinIO 对象路径(无 http(s) 前缀)。
// 模型网关把结果转存 OSS 后返回该裸路径,消费方据此识别"已是文件路径"而不再重复上传。
// 对象命名规则见 oss/minio 的 ensureBucketAndObjectName,格式变化只需改这一处。
func IsOSSPath(s string) bool {
return ossObjectPathPattern.MatchString(s)
}
// GetFileAddressPrefix 拼接图片前缀地址
func GetFileAddressPrefix(ctx context.Context) (imageUrl string, err error) {
// 拼接图片前缀地址
bucketName, err := GetBucketName(ctx)
if err != nil {
return
}
imageUrl = fmt.Sprintf("%s/%s", g.Cfg().MustGet(ctx, "filePrefix").String(), bucketName)
return
}
// GetBucketName 获取bucket名称
func GetBucketName(ctx context.Context) (bucketName string, err error) {
user, err := GetUserInfo(ctx)
if err != nil {
return
}
bucketName = fmt.Sprintf("tenantid-%d", user.TenantId)
return
}
+132 -22
View File
@@ -13,7 +13,7 @@ import (
"sync/atomic" "sync/atomic"
"time" "time"
"gitea.com/red-future/common/beans" "gitea.redpowerfuture.com/red-future/common/beans"
"github.com/gogf/gf/v2/container/gvar" "github.com/gogf/gf/v2/container/gvar"
"github.com/gogf/gf/v2/database/gredis" "github.com/gogf/gf/v2/database/gredis"
"github.com/gogf/gf/v2/errors/gcode" "github.com/gogf/gf/v2/errors/gcode"
@@ -389,27 +389,6 @@ func intPow10(n int) int {
return result return result
} }
// GetFileAddressPrefix 拼接图片前缀地址
func GetFileAddressPrefix(ctx context.Context) (imageUrl string, err error) {
// 拼接图片前缀地址
bucketName, err := GetBucketName(ctx)
if err != nil {
return
}
imageUrl = fmt.Sprintf("%s/%s", g.Cfg().MustGet(ctx, "filePrefix").String(), bucketName)
return
}
// GetBucketName 获取bucket名称
func GetBucketName(ctx context.Context) (bucketName string, err error) {
user, err := GetUserInfo(ctx)
if err != nil {
return
}
bucketName = fmt.Sprintf("tenantid-%d", user.TenantId)
return
}
// Lock 分布式锁 // Lock 分布式锁
func Lock(ctx context.Context, key string, expireSeconds int64, fn func(ctx context.Context) error) (success bool, err error) { func Lock(ctx context.Context, key string, expireSeconds int64, fn func(ctx context.Context) error) (success bool, err error) {
limit := 3 limit := 3
@@ -460,3 +439,134 @@ func IsLocalIP(ip string) bool {
} }
return false return false
} }
// GetLocalIP 获取本机有效的局域网 IPv4 地址
func GetLocalIP() string {
addrs, err := net.InterfaceAddrs()
if err != nil {
return "127.0.0.1"
}
var validIPs []string
for _, addr := range addrs {
ipnet, ok := addr.(*net.IPNet)
if !ok {
continue
}
ip := ipnet.IP
if isIPValid(ip) {
validIPs = append(validIPs, ip.String())
}
}
// 优先返回非 169.254.x.x 的 IP
for _, ip := range validIPs {
if !strings.HasPrefix(ip, "169.254.") {
return ip
}
}
// 其次返回 169.254.x.x(最后的选择)
if len(validIPs) > 0 {
return validIPs[0]
}
return "127.0.0.1"
}
// isIPValid 判断 IP 是否有效
func isIPValid(ip net.IP) bool {
// 不是 loopback (127.0.0.1)
if ip.IsLoopback() {
return false
}
// 是 IPv4
if ip.To4() == nil {
return false
}
// 不是链路本地地址 (169.254.0.0/16)
if ip[0] == 169 && ip[1] == 254 {
return false
}
// 不是组播地址
if ip.IsMulticast() {
return false
}
// 不是未指定地址 (0.0.0.0)
if ip.IsUnspecified() {
return false
}
return true
}
func GetServerPort(ctx context.Context) string {
address := g.Cfg().MustGet(ctx, "server.address", ":8080").String()
// address 格式如 ":3009",去掉冒号
if strings.HasPrefix(address, ":") {
return address[1:]
}
return "8080"
}
// GetLocalAddress 获取局域网地址(IP:端口)
func GetLocalAddress(ctx context.Context) string {
ip := GetLocalIP()
port := GetServerPort(ctx)
if port == "80" || port == "443" {
return ip
}
return ip + ":" + port
}
// GetSchemaFromRequest 从当前请求中获取协议(http/https)
func GetSchemaFromRequest(ctx context.Context) string {
r := g.RequestFromCtx(ctx)
if r == nil {
return "http"
}
// 1. 代理场景:X-Forwarded-Proto
if proto := r.Header.Get("X-Forwarded-Proto"); proto != "" {
return proto
}
// 2. 代理场景:X-Forwarded-Scheme
if proto := r.Header.Get("X-Forwarded-Scheme"); proto != "" {
return proto
}
// 3. TLS 连接(直接 HTTPS
if r.TLS != nil {
return "https"
}
// 4. 默认 HTTP(这行很重要!)
return "http" // ← 确保有这行
}
// GetLocalBaseURL 获取局域网基础 URL(动态协议 + IP + 端口)
func GetLocalBaseURL(ctx context.Context) string {
schema := GetSchemaFromRequest(ctx)
addr := GetLocalAddress(ctx)
return schema + "://" + addr
}
// GetCallbackURL 获取回调地址(完整 URL)
func GetCallbackURL(ctx context.Context, path string) string {
//baseURL := GetLocalBaseURL(ctx)
baseURL := "http://" + GetLocalAddress(ctx)
// 确保 path 以 / 开头
if !strings.HasPrefix(path, "/") {
path = "/" + path
}
return baseURL + path
}
+75
View File
@@ -0,0 +1,75 @@
package websocket
import (
"context"
"sync"
"sync/atomic"
"time"
"github.com/gogf/gf/v2/encoding/gjson"
"github.com/gogf/gf/v2/os/glog"
"github.com/gorilla/websocket"
)
// WsConnection 单个WebSocket连接,Metadata 存放业务自定义数据
type WsConnection struct {
SessionId string
Conn *websocket.Conn
Headers map[string]string
Metadata sync.Map // 业务数据:FlowId, execCancel 等
writeMu sync.Mutex // 保护 websocket.Conn 并发写(WriteMessage + WriteControl
closeCancel context.CancelFunc
closed int32
}
// SetMeta 设置业务元数据
func (c *WsConnection) SetMeta(key string, value interface{}) {
c.Metadata.Store(key, value)
}
// GetMeta 获取业务元数据
func (c *WsConnection) GetMeta(key string) (interface{}, bool) {
return c.Metadata.Load(key)
}
// GetMetaT 泛型版 GetMeta,省去外部类型断言
func GetMetaT[T any](c *WsConnection, key string) (T, bool) {
val, ok := c.Metadata.Load(key)
if !ok {
var zero T
return zero, false
}
t, ok := val.(T)
return t, ok
}
// IsClosed 连接是否已关闭
func (c *WsConnection) IsClosed() bool {
return atomic.LoadInt32(&c.closed) == 1
}
// WriteControl 带写锁保护的 WriteControl,用于心跳 Ping / Pong / Close帧
func (c *WsConnection) WriteControl(msgType int, data []byte, deadline time.Time) error {
c.writeMu.Lock()
defer c.writeMu.Unlock()
_ = c.Conn.SetWriteDeadline(deadline)
return c.Conn.WriteControl(msgType, data, deadline)
}
// WriteJSON 业务层外部写入入口,共享 writeMu 与心跳/Pong 互斥
func (c *WsConnection) WriteJSON(data interface{}) error {
jsonBytes, err := gjson.Encode(data)
if err != nil {
glog.Errorf(context.Background(), "json encode failed: %v", err)
return err
}
c.writeMu.Lock()
_ = c.Conn.SetWriteDeadline(time.Now().Add(30 * time.Second))
err = c.Conn.WriteMessage(websocket.TextMessage, jsonBytes)
c.writeMu.Unlock()
if err != nil {
glog.Debugf(context.Background(), "websocket write failed: %v", err)
}
return err
}
+12
View File
@@ -0,0 +1,12 @@
package websocket
import "time"
const (
DefaultReadTimeout = 90 * time.Second
DefaultWriteTimeout = 10 * time.Second
DefaultHeartbeatInterval = 30 * time.Second
DefaultWorkerPoolSize = 50
DefaultMaxConnections = 2000
DefaultConnKeyPrefix = "ws:"
)
+72
View File
@@ -0,0 +1,72 @@
package websocket
import (
"context"
netHttp "net/http"
"time"
)
// MessageHandler 业务消息处理函数
type MessageHandler func(ctx context.Context, conn *WsConnection, payload interface{})
// WsMessage 通用入站消息
type WsMessage struct {
Type string `json:"type"`
Payload interface{} `json:"payload,omitempty"`
}
// WsPushMsg 通用出站推送消息(业务可自行扩展字段)
type WsPushMsg struct {
Type string `json:"type"`
Message string `json:"message,omitempty"`
Data interface{} `json:"data,omitempty"`
Error string `json:"error,omitempty"`
}
// ServerOptions 服务配置
type ServerOptions struct {
readTimeout time.Duration
writeTimeout time.Duration
heartbeatInterval time.Duration
workerPoolSize int
maxConnections int
connKeyPrefix string
checkOrigin func(r *netHttp.Request) bool
}
// ServerOption 配置函数
type ServerOption func(*ServerOptions)
func WithReadTimeout(d time.Duration) ServerOption {
return func(o *ServerOptions) { o.readTimeout = d }
}
func WithWriteTimeout(d time.Duration) ServerOption {
return func(o *ServerOptions) { o.writeTimeout = d }
}
func WithHeartbeatInterval(d time.Duration) ServerOption {
return func(o *ServerOptions) { o.heartbeatInterval = d }
}
func WithWorkerPoolSize(n int) ServerOption {
return func(o *ServerOptions) { o.workerPoolSize = n }
}
func WithMaxConnections(n int) ServerOption {
return func(o *ServerOptions) { o.maxConnections = n }
}
func WithConnKeyPrefix(p string) ServerOption {
return func(o *ServerOptions) { o.connKeyPrefix = p }
}
func WithCheckOrigin(fn func(r *netHttp.Request) bool) ServerOption {
return func(o *ServerOptions) { o.checkOrigin = fn }
}
func defaultOptions() ServerOptions {
return ServerOptions{
readTimeout: DefaultReadTimeout,
writeTimeout: DefaultWriteTimeout,
heartbeatInterval: DefaultHeartbeatInterval,
workerPoolSize: DefaultWorkerPoolSize,
maxConnections: DefaultMaxConnections,
connKeyPrefix: DefaultConnKeyPrefix,
checkOrigin: func(r *netHttp.Request) bool { return true },
}
}
+326
View File
@@ -0,0 +1,326 @@
package websocket
import (
"context"
"errors"
"fmt"
"sync"
"sync/atomic"
"time"
"github.com/gogf/gf/v2/container/gmap"
"github.com/gogf/gf/v2/encoding/gjson"
"github.com/gogf/gf/v2/frame/g"
"github.com/gogf/gf/v2/net/ghttp"
"github.com/gogf/gf/v2/os/glog"
"github.com/gogf/gf/v2/os/grpool"
"github.com/google/uuid"
"github.com/gorilla/websocket"
)
// WsServer 泛化 WebSocket 服务器
type WsServer struct {
connections *gmap.StrAnyMap
upgrader websocket.Upgrader
workerPool *grpool.Pool
handlers map[string]MessageHandler
handlerMu sync.RWMutex
opts ServerOptions
svcClosed int32
closeOnce sync.Once
}
// NewWsServer 创建泛化 WebSocket 服务器
func NewWsServer(opts ...ServerOption) *WsServer {
o := defaultOptions()
for _, opt := range opts {
opt(&o)
}
return &WsServer{
connections: gmap.NewStrAnyMap(true),
upgrader: websocket.Upgrader{
ReadBufferSize: 1024,
WriteBufferSize: 1024,
CheckOrigin: o.checkOrigin,
},
workerPool: grpool.New(o.workerPoolSize),
handlers: make(map[string]MessageHandler),
opts: o,
svcClosed: 0,
}
}
// OnMessage 注册业务消息处理器
func (s *WsServer) OnMessage(msgType string, handler MessageHandler) {
s.handlerMu.Lock()
defer s.handlerMu.Unlock()
s.handlers[msgType] = handler
}
// Upgrade 将 HTTP 连接升级为 WebSocket 并注册到连接池
func (s *WsServer) Upgrade(ctx context.Context, r *ghttp.Request, sessionId string) (*WsConnection, error) {
if g.IsEmpty(sessionId) {
sessionId = uuid.NewString()
}
if atomic.LoadInt32(&s.svcClosed) == 1 {
return nil, errors.New("websocket server is closed")
}
if s.connections.Size() >= s.opts.maxConnections {
return nil, errors.New("too many online websocket connections")
}
wsConn, err := s.upgrader.Upgrade(r.Response.Writer, r.Request, nil)
if err != nil {
return nil, fmt.Errorf("upgrade failed: %w", err)
}
headers := make(map[string]string)
for k, v := range r.Request.Header {
if len(v) > 0 {
headers[k] = v[0]
}
}
key := s.opts.connKeyPrefix + sessionId
// 踢下线旧连接
s.kickOld(key)
baseCtx := context.WithoutCancel(ctx)
closeCtx, closeCancel := context.WithCancel(baseCtx)
wc := &WsConnection{
SessionId: sessionId,
Conn: wsConn,
Headers: headers,
closeCancel: closeCancel,
closed: 0,
}
s.connections.Set(key, wc)
// 连接成功回执
_ = s.writeJSON(closeCtx, wc, &WsPushMsg{Type: "ack", Message: "WebSocket连接成功", Data: map[string]any{
"sessionId": sessionId,
}})
// Pong 心跳回调,重置读超时
wsConn.SetPongHandler(func(string) error {
_ = wsConn.SetReadDeadline(time.Now().Add(s.opts.readTimeout))
return nil
})
go s.handleConnection(closeCtx, key, wc)
return wc, nil
}
// PushToSession 向指定会话推送消息
func (s *WsServer) PushToSession(ctx context.Context, sessionId string, msg *WsPushMsg) {
key := s.opts.connKeyPrefix + sessionId
val := s.connections.Get(key)
if val == nil {
return
}
wc, ok := val.(*WsConnection)
if !ok || wc.IsClosed() {
return
}
_ = s.writeJSON(ctx, wc, msg)
}
// GetOnlineSessions 获取在线会话列表
func (s *WsServer) GetOnlineSessions() []string {
var sessions []string
prefixLen := len(s.opts.connKeyPrefix)
s.connections.Iterator(func(key string, _ interface{}) bool {
if len(key) > prefixLen {
sessions = append(sessions, key[prefixLen:])
}
return true
})
return sessions
}
// Close 全局优雅关闭
func (s *WsServer) Close() {
s.closeOnce.Do(func() {
atomic.StoreInt32(&s.svcClosed, 1)
s.workerPool.Close()
s.connections.LockFunc(func(m map[string]interface{}) {
for _, val := range m {
wc, ok := val.(*WsConnection)
if !ok {
continue
}
if atomic.CompareAndSwapInt32(&wc.closed, 0, 1) {
if wc.closeCancel != nil {
wc.closeCancel()
}
_ = wc.Conn.Close()
}
}
})
s.connections.Clear()
})
}
// ====================== 内部方法 ======================
// kickOld 踢掉同session旧连接,不再主动remove,由旧连接defer清理
func (s *WsServer) kickOld(key string) {
val := s.connections.Get(key)
if val == nil {
return
}
old, ok := val.(*WsConnection)
if !ok {
return
}
if atomic.CompareAndSwapInt32(&old.closed, 0, 1) {
if old.closeCancel != nil {
old.closeCancel()
}
_ = old.Conn.Close()
}
}
// heartbeatLoop 心跳发送协程,入参改为 *WsConnection,复用写锁
func (s *WsServer) heartbeatLoop(ctx context.Context, wc *WsConnection, done <-chan struct{}) {
ticker := time.NewTicker(s.opts.heartbeatInterval)
defer ticker.Stop()
conn := wc.Conn
for {
select {
case <-ticker.C:
wc.writeMu.Lock()
_ = conn.SetWriteDeadline(time.Now().Add(s.opts.writeTimeout))
err := conn.WriteControl(websocket.PingMessage, nil, time.Now().Add(s.opts.writeTimeout))
wc.writeMu.Unlock()
if err != nil {
glog.Debugf(ctx, "heartbeat ping failed: %v", err)
return
}
case <-done:
return
case <-ctx.Done():
return
}
}
}
func (s *WsServer) handleConnection(ctx context.Context, key string, wc *WsConnection) {
conn := wc.Conn
defer func() {
if atomic.CompareAndSwapInt32(&wc.closed, 0, 1) {
if wc.closeCancel != nil {
wc.closeCancel()
}
_ = conn.Close()
}
// 关键修复:只删除自身实例,防止旧连接误删新连接
s.connections.LockFunc(func(m map[string]interface{}) {
if v, exist := m[key]; exist && v == wc {
delete(m, key)
}
})
}()
done := make(chan struct{})
defer close(done)
go s.heartbeatLoop(ctx, wc, done)
_ = conn.SetReadDeadline(time.Now().Add(s.opts.readTimeout))
for {
select {
case <-ctx.Done():
return
default:
}
msgType, data, err := conn.ReadMessage()
if err != nil {
// 正常关闭不打error日志
if !websocket.IsUnexpectedCloseError(err,
websocket.CloseNormalClosure,
websocket.CloseGoingAway,
websocket.CloseNoStatusReceived,
) {
glog.Debugf(ctx, "normal close: %s, err: %v", key, err)
} else {
glog.Infof(ctx, "unexpected close: %s, err: %v", key, err)
}
break
}
_ = conn.SetReadDeadline(time.Now().Add(s.opts.readTimeout))
switch msgType {
case websocket.PingMessage:
wc.writeMu.Lock()
_ = conn.SetWriteDeadline(time.Now().Add(s.opts.writeTimeout))
_ = conn.WriteMessage(websocket.PongMessage, nil)
wc.writeMu.Unlock()
continue
case websocket.CloseMessage:
return
case websocket.BinaryMessage, websocket.TextMessage:
default:
continue
}
if len(data) == 0 {
continue
}
var msg WsMessage
if err := gjson.Unmarshal(data, &msg); err != nil {
_ = s.writeJSON(ctx, wc, &WsPushMsg{Type: "error", Message: "消息格式错误", Error: err.Error()})
continue
}
s.handlerMu.RLock()
handler, exists := s.handlers[msg.Type]
s.handlerMu.RUnlock()
if !exists {
_ = s.writeJSON(ctx, wc, &WsPushMsg{Type: "error", Message: fmt.Sprintf("未知消息类型: %s", msg.Type)})
continue
}
// 【重要修复】投递到workerPool,避免业务阻塞读循环
taskCtx := ctx
payload := msg.Payload
if err := s.workerPool.Add(taskCtx, func(ctx context.Context) {
handler(ctx, wc, payload)
}); err != nil {
_ = s.writeJSON(ctx, wc, &WsPushMsg{
Type: "error",
Message: "服务繁忙,任务队列已满",
})
}
}
}
// writeJSON 统一写入消息,入参改为 *WsConnection,带并发写锁
func (s *WsServer) writeJSON(ctx context.Context, wc *WsConnection, data interface{}) error {
wc.writeMu.Lock()
defer wc.writeMu.Unlock()
jsonBytes, err := gjson.Encode(data)
if err != nil {
glog.Errorf(ctx, "json encode failed: %v", err)
return err
}
_ = wc.Conn.SetWriteDeadline(time.Now().Add(s.opts.writeTimeout))
if err = wc.Conn.WriteMessage(websocket.TextMessage, jsonBytes); err != nil {
glog.Debugf(ctx, "websocket write failed: %v", err)
return err
}
return nil
}