diff --git a/api/postal.proto b/api/postal.proto index d6c410e..5071aae 100644 --- a/api/postal.proto +++ b/api/postal.proto @@ -25,7 +25,7 @@ message Message { message ReqDeliver { string receiver = 1; - Message msg = 2; + Message message = 2; // bool sync = 7; // 是否同步阻塞等待投递结果, 默认false立即返回放到队列消费投递 } diff --git a/cmd/auth/main.go b/cmd/auth/main.go index 06ab7d0..e3a82b6 100644 --- a/cmd/auth/main.go +++ b/cmd/auth/main.go @@ -1 +1,12 @@ package main + +import ( + "encoding/binary" + "fmt" +) + +func main() { + buf := make([]byte, 4) + binary.BigEndian.PutUint32(buf, 12) + fmt.Printf("%v\n", buf) +} diff --git a/cmd/gateway_ws/main.go b/cmd/gateway_ws/main.go index 4af2e58..fdad1c8 100644 --- a/cmd/gateway_ws/main.go +++ b/cmd/gateway_ws/main.go @@ -2,10 +2,13 @@ package main import ( clientv3 "go.etcd.io/etcd/client/v3" + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" "google.golang.org/grpc/grpclog" "sonet/internal/postal/logic" "sonet/pkg/config" "sonet/pkg/grpc/discovery" + "sonet/pkg/grpc/generic" "sonet/pkg/utils/logger" "sonet/pkg/utils/shutdown" ) @@ -15,16 +18,22 @@ func main() { grpclog.SetLoggerV2(logger.Logger) conf := config.LoadConfig(nil, "cmd/gateway_ws") - postalServer := logic.NewPostalServer() // registry and run... etcdClient, err := clientv3.New(conf.Etcd) if err != nil { panic(err) } + + resolver := discovery.NewResolver(etcdClient) + + generic.NewGpcGenericClientFactory(resolver, grpc.WithTransportCredentials(insecure.NewCredentials())) + //server.NewConnHandler() + //server.NewHttpServer() + postalRegistry := discovery.NewRegister(etcdClient) shutdown.AddShutdownHook(postalRegistry.Stop) - + postalServer := logic.NewPostalServer() go func() { err = postalServer.Run(conf.Grpc, postalRegistry) if err != nil { diff --git a/generate.go b/generate.go index ce37340..6fa3169 100644 --- a/generate.go +++ b/generate.go @@ -1,7 +1,25 @@ -package sonet +package main +import "os" + +//go:generate go run generate.go //go:generate protoc --go_out=./api/gen --go-grpc_out=./api/gen ./api/*.proto //go:generate pbjs -t static-module -w es6 -o api/genjs/postal.js api/postal.proto --no-service //go:generate pbjs -t static-module -w es6 -o api/genjs/auth.js api/auth.proto --no-service //go:generate pbjs -t static-module -w es6 -o api/genjs/chat.js api/chat.proto --no-service //go:generate pbjs -t static-module -w es6 -o api/genjs/mahjong.js api/mahjong.proto --no-service + +func main() { + beforeGenerate() +} + +func beforeGenerate() { + err := os.MkdirAll("api/gen", os.ModePerm) + if err != nil { + panic(err) + } + err = os.MkdirAll("api/genjs", os.ModePerm) + if err != nil { + panic(err) + } +} diff --git a/go.mod b/go.mod index b2454bd..d6d100f 100644 --- a/go.mod +++ b/go.mod @@ -18,20 +18,36 @@ require ( require ( github.com/bufbuild/protocompile v0.7.1 // indirect + github.com/bytedance/sonic v1.9.1 // indirect github.com/cespare/xxhash/v2 v2.2.0 // indirect + github.com/chenzhuoyu/base64x v0.0.0-20221115062448-fe3a3abad311 // indirect github.com/coreos/go-semver v0.3.0 // indirect github.com/coreos/go-systemd/v22 v22.3.2 // indirect github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f // indirect github.com/dsnet/golib/unitconv v1.0.2 // indirect github.com/fsnotify/fsnotify v1.7.0 // indirect + github.com/gabriel-vasile/mimetype v1.4.2 // indirect + github.com/gin-contrib/sse v0.1.0 // indirect + github.com/gin-gonic/gin v1.9.1 // indirect + github.com/go-playground/locales v0.14.1 // indirect + github.com/go-playground/universal-translator v0.18.1 // indirect + github.com/go-playground/validator/v10 v10.14.0 // indirect github.com/go-sql-driver/mysql v1.7.0 // indirect + github.com/goccy/go-json v0.10.2 // indirect github.com/gogo/protobuf v1.3.2 // indirect + github.com/gorilla/websocket v1.5.1 // indirect github.com/hashicorp/hcl v1.0.0 // indirect github.com/jinzhu/inflection v1.0.0 // indirect github.com/jinzhu/now v1.1.5 // indirect + github.com/json-iterator/go v1.1.12 // indirect github.com/klauspost/compress v1.17.0 // indirect + github.com/klauspost/cpuid/v2 v2.2.4 // indirect + github.com/leodido/go-urn v1.2.4 // indirect github.com/magiconair/properties v1.8.7 // indirect + github.com/mattn/go-isatty v0.0.19 // indirect github.com/mitchellh/mapstructure v1.5.0 // indirect + github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect + github.com/modern-go/reflect2 v1.0.2 // indirect github.com/nats-io/nkeys v0.4.6 // indirect github.com/nats-io/nuid v1.0.1 // indirect github.com/pelletier/go-toml/v2 v2.1.0 // indirect @@ -42,11 +58,14 @@ require ( github.com/spf13/cast v1.6.0 // indirect github.com/spf13/pflag v1.0.5 // indirect github.com/subosito/gotenv v1.6.0 // indirect + github.com/twitchyliquid64/golang-asm v0.15.1 // indirect + github.com/ugorji/go/codec v1.2.11 // indirect go.etcd.io/etcd/api/v3 v3.5.11 // indirect go.etcd.io/etcd/client/pkg/v3 v3.5.11 // indirect go.uber.org/atomic v1.9.0 // indirect go.uber.org/multierr v1.9.0 // indirect go.uber.org/zap v1.21.0 // indirect + golang.org/x/arch v0.3.0 // indirect golang.org/x/crypto v0.16.0 // indirect golang.org/x/exp v0.0.0-20230905200255-921286631fa9 // indirect golang.org/x/net v0.19.0 // indirect diff --git a/go.sum b/go.sum index 693378c..016feda 100644 --- a/go.sum +++ b/go.sum @@ -4,8 +4,14 @@ github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs= github.com/bsm/gomega v1.27.10 h1:yeMWxP2pV2fG3FgAODIY8EiRE3dy0aeFYt4l7wh6yKA= github.com/bufbuild/protocompile v0.7.1 h1:Kd8fb6EshOHXNNRtYAmLAwy/PotlyFoN0iMbuwGNh0M= github.com/bufbuild/protocompile v0.7.1/go.mod h1:+Etjg4guZoAqzVk2czwEQP12yaxLJ8DxuqCJ9qHdH94= +github.com/bytedance/sonic v1.5.0/go.mod h1:ED5hyg4y6t3/9Ku1R6dU/4KyJ48DZ4jPhfY1O2AihPM= +github.com/bytedance/sonic v1.9.1 h1:6iJ6NqdoxCDr6mbY8h18oSO+cShGSMRGCEo7F2h0x8s= +github.com/bytedance/sonic v1.9.1/go.mod h1:i736AoUSYt75HyZLoJW9ERYxcy6eaN6h4BZXU064P/U= github.com/cespare/xxhash/v2 v2.2.0 h1:DC2CZ1Ep5Y4k3ZQ899DldepgrayRUGE6BBZ/cd9Cj44= github.com/cespare/xxhash/v2 v2.2.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/chenzhuoyu/base64x v0.0.0-20211019084208-fb5309c8db06/go.mod h1:DH46F32mSOjUmXrMHnKwZdA8wcEefY7UVqBKYGjpdQY= +github.com/chenzhuoyu/base64x v0.0.0-20221115062448-fe3a3abad311 h1:qSGYFH7+jGhDF8vLC+iwCD4WpbV1EBDSzWkJODFLams= +github.com/chenzhuoyu/base64x v0.0.0-20221115062448-fe3a3abad311/go.mod h1:b583jCggY9gE99b6G5LEC39OIiVsWj+R97kbl5odCEk= github.com/coreos/go-semver v0.3.0 h1:wkHLiw0WNATZnSG7epLsujiMCgPAc9xhjJ4tgnAxmfM= github.com/coreos/go-semver v0.3.0/go.mod h1:nnelYz7RCh+5ahJtPPxZlU+153eP4D4r3EedlOD2RNk= github.com/coreos/go-systemd/v22 v22.3.2 h1:D9/bQk5vlXQFZ6Kwuu6zaiXJ9oTPe68++AzAJc1DzSI= @@ -20,8 +26,22 @@ github.com/dsnet/golib/unitconv v1.0.2/go.mod h1:86KTUtTJFLreKjc4sS9xE0rhj4lR44O github.com/frankban/quicktest v1.14.6 h1:7Xjx+VpznH+oBnejlPUj8oUpdxnVs4f8XU8WnHkI4W8= github.com/fsnotify/fsnotify v1.7.0 h1:8JEhPFa5W2WU7YfeZzPNqzMP6Lwt7L2715Ggo0nosvA= github.com/fsnotify/fsnotify v1.7.0/go.mod h1:40Bi/Hjc2AVfZrqy+aj+yEI+/bRxZnMJyTJwOpGvigM= +github.com/gabriel-vasile/mimetype v1.4.2 h1:w5qFW6JKBz9Y393Y4q372O9A7cUSequkh1Q7OhCmWKU= +github.com/gabriel-vasile/mimetype v1.4.2/go.mod h1:zApsH/mKG4w07erKIaJPFiX0Tsq9BFQgN3qGY5GnNgA= +github.com/gin-contrib/sse v0.1.0 h1:Y/yl/+YNO8GZSjAhjMsSuLt29uWRFHdHYUb5lYOV9qE= +github.com/gin-contrib/sse v0.1.0/go.mod h1:RHrZQHXnP2xjPF+u1gW/2HnVO7nvIa9PG3Gm+fLHvGI= +github.com/gin-gonic/gin v1.9.1 h1:4idEAncQnU5cB7BeOkPtxjfCSye0AAm1R0RVIqJ+Jmg= +github.com/gin-gonic/gin v1.9.1/go.mod h1:hPrL7YrpYKXt5YId3A/Tnip5kqbEAP+KLuI3SUcPTeU= +github.com/go-playground/locales v0.14.1 h1:EWaQ/wswjilfKLTECiXz7Rh+3BjFhfDFKv/oXslEjJA= +github.com/go-playground/locales v0.14.1/go.mod h1:hxrqLVvrK65+Rwrd5Fc6F2O76J/NuW9t0sjnWqG1slY= +github.com/go-playground/universal-translator v0.18.1 h1:Bcnm0ZwsGyWbCzImXv+pAJnYK9S473LQFuzCbDbfSFY= +github.com/go-playground/universal-translator v0.18.1/go.mod h1:xekY+UJKNuX9WP91TpwSH2VMlDf28Uj24BCp08ZFTUY= +github.com/go-playground/validator/v10 v10.14.0 h1:vgvQWe3XCz3gIeFDm/HnTIbj6UGmg/+t63MyGU2n5js= +github.com/go-playground/validator/v10 v10.14.0/go.mod h1:9iXMNT7sEkjXb0I+enO7QXmzG6QCsPWY4zveKFVRSyU= github.com/go-sql-driver/mysql v1.7.0 h1:ueSltNNllEqE3qcWBTD0iQd3IpL/6U+mJxLkazJ7YPc= github.com/go-sql-driver/mysql v1.7.0/go.mod h1:OXbVy3sEdcQ2Doequ6Z5BW6fXNQTmx+9S1MCJN5yJMI= +github.com/goccy/go-json v0.10.2 h1:CrxCmQqYDkv1z7lO7Wbh2HN93uovUHgrECaO5ZrCXAU= +github.com/goccy/go-json v0.10.2/go.mod h1:6MelG93GURQebXPDq3khkgXZkazVtN9CRI+MGFi0w8I= github.com/godbus/dbus/v5 v5.0.4/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA= github.com/gogo/protobuf v1.3.2 h1:Ov1cvc58UF3b5XjBnZv7+opcTcQFZebYjWzi34vdm4Q= github.com/gogo/protobuf v1.3.2/go.mod h1:P1XiOD3dCwIKUDQYPy72D8LYyHL2YPYrpS2s69NZV8Q= @@ -31,6 +51,9 @@ github.com/golang/protobuf v1.5.3/go.mod h1:XVQd3VNwM+JqD3oG2Ue2ip4fOMUkwXdXDdiu github.com/google/go-cmp v0.3.0/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU= github.com/google/go-cmp v0.5.5/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= +github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= +github.com/gorilla/websocket v1.5.1 h1:gmztn0JnHVt9JZquRuzLw3g4wouNVzKL15iLr/zn/QY= +github.com/gorilla/websocket v1.5.1/go.mod h1:x3kM2JMyaluk02fnUJpQuwD2dCS5NDG2ZHL0uE0tcaY= github.com/hashicorp/hcl v1.0.0 h1:0Anlzjpi4vEasTeNFn2mLJgTSwt0+6sfsiTG8qcWGx4= github.com/hashicorp/hcl v1.0.0/go.mod h1:E5yfLk+7swimpb2L/Alb/PJmXilQ/rhwaUYs4T20WEQ= github.com/jhump/protoreflect v1.15.4 h1:mrwJhfQGGljwvR/jPEocli8KA6G9afbQpH8NY2wORcI= @@ -39,19 +62,33 @@ github.com/jinzhu/inflection v1.0.0 h1:K317FqzuhWc8YvSVlFMCCUb36O/S9MCKRDI7QkRKD github.com/jinzhu/inflection v1.0.0/go.mod h1:h+uFLlag+Qp1Va5pdKtLDYj+kHp5pxUVkryuEj+Srlc= github.com/jinzhu/now v1.1.5 h1:/o9tlHleP7gOFmsnYNz3RGnqzefHA47wQpKrrdTIwXQ= github.com/jinzhu/now v1.1.5/go.mod h1:d3SSVoowX0Lcu0IBviAWJpolVfI5UJVZZ7cO71lE/z8= +github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM= +github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo= github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI2bnpBCr8= github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck= github.com/klauspost/compress v1.17.0 h1:Rnbp4K9EjcDuVuHtd0dgA4qNuv9yKDYKK1ulpJwgrqM= github.com/klauspost/compress v1.17.0/go.mod h1:ntbaceVETuRiXiv4DpjP66DpAtAGkEQskQzEyD//IeE= +github.com/klauspost/cpuid/v2 v2.0.9/go.mod h1:FInQzS24/EEf25PyTYn52gqo7WaD8xa0213Md/qVLRg= +github.com/klauspost/cpuid/v2 v2.2.4 h1:acbojRNwl3o09bUq+yDCtZFc1aiwaAAxtcn8YkZXnvk= +github.com/klauspost/cpuid/v2 v2.2.4/go.mod h1:RVVoqg1df56z8g3pUjL/3lE5UfnlrJX8tyFgg4nqhuY= github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ= github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/leodido/go-urn v1.2.4 h1:XlAE/cm/ms7TE/VMVoduSpNBoyc2dOxHs5MZSwAN63Q= +github.com/leodido/go-urn v1.2.4/go.mod h1:7ZrI8mTSeBSHl/UaRyKQW1qZeMgak41ANeCNaVckg+4= github.com/magiconair/properties v1.8.7 h1:IeQXZAiQcpL9mgcAe1Nu6cX9LLw6ExEHKjN0VQdvPDY= github.com/magiconair/properties v1.8.7/go.mod h1:Dhd985XPs7jluiymwWYZ0G4Z61jb3vdS329zhj2hYo0= +github.com/mattn/go-isatty v0.0.19 h1:JITubQf0MOLdlGRuRq+jtsDlekdYPia9ZFsB8h/APPA= +github.com/mattn/go-isatty v0.0.19/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= github.com/mitchellh/mapstructure v1.5.0 h1:jeMsZIYE/09sWLaz43PL7Gy6RuMjD2eJVyuac5Z2hdY= github.com/mitchellh/mapstructure v1.5.0/go.mod h1:bFUtVrKA4DC2yAKiSyO/QUcy7e+RRV2QTWOzhPopBRo= +github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= +github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd h1:TRLaZ9cD/w8PVh93nsPXa1VrQ6jlwL5oN8l14QlcNfg= +github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= +github.com/modern-go/reflect2 v1.0.2 h1:xBagoLtFs94CBntxluKeaWgTMpvLxC4ur3nMaC9Gz0M= +github.com/modern-go/reflect2 v1.0.2/go.mod h1:yWuevngMOJpCy52FWWMvUC8ws7m/LJsjYzDa0/r8luk= github.com/nats-io/nats.go v1.31.0 h1:/WFBHEc/dOKBF6qf1TZhrdEfTmOZ5JzdJ+Y3m6Y/p7E= github.com/nats-io/nats.go v1.31.0/go.mod h1:di3Bm5MLsoB4Bx61CBTsxuarI36WbhAwOm8QrW39+i8= github.com/nats-io/nkeys v0.4.6 h1:IzVe95ru2CT6ta874rt9saQRkWfe2nFj1NtvYSLqMzY= @@ -90,10 +127,16 @@ github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UV github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU= +github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4= +github.com/stretchr/testify v1.8.2/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4= github.com/stretchr/testify v1.8.4 h1:CcVxjf3Q8PM0mHUKJCdn+eZZtm5yQwehR5yeSVQQcUk= github.com/stretchr/testify v1.8.4/go.mod h1:sz/lmYIOXD/1dqDmKjjqLyZ2RngseejIcXlSw2iwfAo= github.com/subosito/gotenv v1.6.0 h1:9NlTDc1FTs4qu0DDq7AEtTPNw6SVm7uBMsUCUjABIf8= github.com/subosito/gotenv v1.6.0/go.mod h1:Dk4QP5c2W3ibzajGcXpNraDfq2IrhjMIvMSWPKKo0FU= +github.com/twitchyliquid64/golang-asm v0.15.1 h1:SU5vSMR7hnwNxj24w34ZyCi/FmDZTkS4MhqMhdFk5YI= +github.com/twitchyliquid64/golang-asm v0.15.1/go.mod h1:a1lVb/DtPvCB8fslRZhAngC2+aY1QWCk3Cedj/Gdt08= +github.com/ugorji/go/codec v1.2.11 h1:BMaWp1Bb6fHwEtbplGBGJ498wD+LKlNSl25MjdZY4dU= +github.com/ugorji/go/codec v1.2.11/go.mod h1:UNopzCgEMSXjBc6AOMqYvWC1ktqTAfzJZUZgYf6w6lg= github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= github.com/yuin/goldmark v1.2.1/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= github.com/yuin/goldmark v1.3.5/go.mod h1:mwnBkeHKe2W/ZEtQ+71ViKU8L12m81fl3OWwC1Zlc8k= @@ -113,6 +156,9 @@ go.uber.org/multierr v1.9.0 h1:7fIwc/ZtS0q++VgcfqFDxSBZVv/Xo49/SYnDFupUwlI= go.uber.org/multierr v1.9.0/go.mod h1:X2jQV1h+kxSjClGpnseKVIxpmcjrj7MNnI0bnlfKTVQ= go.uber.org/zap v1.21.0 h1:WefMeulhovoZ2sYXz7st6K0sLj7bBhpiFaud4r4zST8= go.uber.org/zap v1.21.0/go.mod h1:wjWOCqI0f2ZZrJF/UufIOkiC8ii6tm1iqIsLo76RfJw= +golang.org/x/arch v0.0.0-20210923205945-b76863e36670/go.mod h1:5om86z9Hs0C8fWVUuoMHwpExlXzs5Tkyp9hOrfG7pp8= +golang.org/x/arch v0.3.0 h1:02VY4/ZcO/gBOH6PUaoiptASxtXU10jazRCP865E97k= +golang.org/x/arch v0.3.0/go.mod h1:5om86z9Hs0C8fWVUuoMHwpExlXzs5Tkyp9hOrfG7pp8= golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= @@ -144,7 +190,9 @@ golang.org/x/sys v0.0.0-20200930185726-fdedc70b468f/go.mod h1:h1NjWce9XRLGQEsW7w golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210330210617-4fbd30eecc44/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210510120138-977fb7262007/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.0.0-20220704084225-05e143d24a9e/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20220715151400-c0bba94af5f8/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.15.0 h1:h48lPFYpsTvQJZF4EKyI4aLHaev3CxivZmv7yZig9pc= golang.org/x/sys v0.15.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= @@ -190,3 +238,4 @@ gorm.io/driver/mysql v1.5.2/go.mod h1:pQLhh1Ut/WUAySdTHwBpBv6+JKcj+ua4ZFx1QQTBzb gorm.io/gorm v1.25.2-0.20230530020048-26663ab9bf55/go.mod h1:L4uxeKpfBml98NYqVqwAdmV1a2nBtAec/cf3fpucW/k= gorm.io/gorm v1.25.5 h1:zR9lOiiYf09VNh5Q1gphfyia1JpiClIWG9hQaxB/mls= gorm.io/gorm v1.25.5/go.mod h1:hbnx/Oo0ChWMn1BIhpy1oYozzpM15i4YPuHDmfYtwg8= +rsc.io/pdf v0.1.1/go.mod h1:n8OzWcQ6Sp37PL01nO98y4iUCRdTGarVfzxY20ICaU4= diff --git a/internal/chat/logic/chat_server.go b/internal/chat/logic/chat_server.go index fd85711..b686da9 100644 --- a/internal/chat/logic/chat_server.go +++ b/internal/chat/logic/chat_server.go @@ -4,13 +4,20 @@ import ( "context" "google.golang.org/protobuf/types/known/emptypb" "sonet/api/gen/chat" + "sonet/api/gen/postal" ) type ChatServer struct { chat.UnimplementedChatServer + postalCli postal.PostalClient +} + +func NewChatServer() *ChatServer { + return &ChatServer{} } func (s *ChatServer) Send(ctx context.Context, send *chat.ReqSend) (*chat.ResSend, error) { + return nil, nil } diff --git a/internal/gateway_ws/server/http_server.go b/internal/gateway_ws/server/http_server.go new file mode 100644 index 0000000..666613e --- /dev/null +++ b/internal/gateway_ws/server/http_server.go @@ -0,0 +1,63 @@ +package server + +import ( + "fmt" + "github.com/gin-gonic/gin" + "github.com/gorilla/websocket" + "net/http" + "sonet/pkg/utils/conver" + "sonet/pkg/utils/resp" + "time" +) + +type HttpServer struct { + connHandler *ConnHandler +} + +func NewHttpServer(connHandler *ConnHandler) *HttpServer { + return &HttpServer{ + connHandler: connHandler, + } +} + +func (s *HttpServer) Run(port int) error { + server := gin.Default() + server.GET("/ws", s.upgrade) + return server.Run(fmt.Sprintf(":%d", port)) +} + +var ( + HandshakeTimeout = 3 * time.Second + ReadDeadline = 5 * time.Second + WriteDeadline = 5 * time.Second + // PongWait Time allowed to read the next pong message from the peer. + PongWait = 60 * time.Second + + MaxMessageSize = conver.MustParseDataUnitInt("4M") + ReadBufferSize = conver.MustParseDataUnitInt("4Ki") + WriteBufferSize = conver.MustParseDataUnitInt("4Ki") +) + +func (s *HttpServer) upgrade(ctx *gin.Context) { + upgrader := websocket.Upgrader{ + ReadBufferSize: ReadBufferSize, + WriteBufferSize: WriteBufferSize, + HandshakeTimeout: HandshakeTimeout, + CheckOrigin: func(r *http.Request) bool { + return true + }, + } + conn, err := upgrader.Upgrade(ctx.Writer, ctx.Request, nil) + if err != nil { + ctx.JSON(http.StatusOK, resp.Error(err.Error())) + return + } + + // https://github.com/gorilla/websocket/blob/a68708917c6a4f06314ab4e52493cc61359c9d42/examples/chat/conn.go#L50 + conn.SetReadLimit(int64(MaxMessageSize)) + conn.SetPongHandler(func(string) error { + return conn.SetReadDeadline(time.Now().Add(PongWait)) + }) + + go s.connHandler.handleConn(conn) +} diff --git a/internal/gateway_ws/server/ws_server.go b/internal/gateway_ws/server/ws_server.go new file mode 100644 index 0000000..e07d1d6 --- /dev/null +++ b/internal/gateway_ws/server/ws_server.go @@ -0,0 +1,89 @@ +package server + +import ( + "context" + "fmt" + "github.com/gorilla/websocket" + "runtime/debug" + "sonet/internal/gateway_ws/session" + "sonet/pkg/grpc/generic" + "sonet/pkg/protocol" + "sonet/pkg/utils/logger" +) + +type ConnHandler struct { + grpcFactory *generic.GrpcGenericClientFactory +} + +func NewConnHandler(grpcFactory *generic.GrpcGenericClientFactory) *ConnHandler { + return &ConnHandler{ + grpcFactory: grpcFactory, + } +} + +func (c *ConnHandler) handleConn(conn *websocket.Conn) { + client := session.NewNetClient(conn, ReadDeadline, WriteDeadline) + + defer func() { + // 捕获其他错误 + if r := recover(); r != nil { + logger.Error("NetClient recover error: ", r) + // 输出堆栈信息 + logger.Error("NetClient recover error stack: ", string(debug.Stack())) + } + }() + + defer client.Close() + + for { + message, err := client.ReadMessage() + if err != nil { + // TODO 连接关闭,mq发送关闭事件 + if websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway) { + logger.Error("unexpected close error: ", err) + } else { + // 读失败 + logger.Errorf("NetClient ReadMessage error: %T, %v", err, err) + } + return + } + + payload, err := protocol.Decode(message) + if err != nil { + logger.Errorf("decode message error: len=%d", len(message), err) + return + } + fmt.Printf("%v\n", payload) + + // grpc generic call + ctx := context.Background() + grpcClient, err := c.grpcFactory.GetClient(ctx, payload.Header.Svc) + if err != nil { + logger.Error("get grpc generic client error: ", err) + continue + } + resp, err := grpcClient.InvokeUnary(ctx, payload.Header.Target, payload.Body) + if err != nil { + logger.Error("grpc generic call error: ", err) + continue + } + fmt.Println(resp) + + // write response + payload.Header.Type = protocol.TypeResponse + payload.Header.Svc = "" + payload.Header.Target = "" + payload.Body, err = resp.Marshal() + if err != nil { + logger.Error("generic call response marshal error: ", err) + continue + } + resMessage, err := protocol.Encode(payload) + if err != nil { + logger.Error("grpc generic call error: ", err) + continue + } + client.MustWrite(resMessage) + } + +} diff --git a/internal/gateway_ws/session/net_account.go b/internal/gateway_ws/session/net_account.go new file mode 100644 index 0000000..a081160 --- /dev/null +++ b/internal/gateway_ws/session/net_account.go @@ -0,0 +1,15 @@ +package session + +import ( + "sync" +) + +// NetAccount 已认证的长连接用户 +type NetAccount struct { + Id int64 + UserName string + Avatar string + Client *NetClient + + Lock *sync.Mutex +} diff --git a/internal/gateway_ws/session/net_client.go b/internal/gateway_ws/session/net_client.go new file mode 100644 index 0000000..647ac61 --- /dev/null +++ b/internal/gateway_ws/session/net_client.go @@ -0,0 +1,77 @@ +package session + +import ( + "fmt" + "github.com/gorilla/websocket" + "sonet/pkg/utils/logger" + "sync" + "sync/atomic" + "time" +) + +// NetClient 长连接客户端 +type NetClient struct { + Conn *websocket.Conn + writeDeadline, readDeadline time.Duration + Online *atomic.Bool + Account *NetAccount + WriteLock *sync.Mutex +} + +func NewNetClient(conn *websocket.Conn, readDeadline, writeDeadline time.Duration) *NetClient { + online := &atomic.Bool{} + online.Store(true) + + return &NetClient{ + Conn: conn, + readDeadline: readDeadline, + writeDeadline: writeDeadline, + WriteLock: new(sync.Mutex), + Online: online, + } +} + +func (c *NetClient) Close() { + err := c.Conn.Close() + if err != nil { + logger.Error("NetClient Close error: ", err) + } +} + +func (c *NetClient) Write(bytes []byte) (err error) { + c.WriteLock.Lock() + defer c.WriteLock.Unlock() + err = c.Conn.SetWriteDeadline(time.Now().Add(10 * time.Second)) + if err != nil { + err = fmt.Errorf("NetClient SetWriteDeadline error: %s", err.Error()) + return + } + err = c.Conn.WriteMessage(websocket.BinaryMessage, bytes) + if err != nil { + err = fmt.Errorf("NetClient write message error: %s", err.Error()) + } + return +} + +func (c *NetClient) MustWrite(bytes []byte) { + err := c.Write(bytes) + if err != nil { + logger.Error(err) + return + } +} + +func (c *NetClient) ReadMessage() (bytes []byte, err error) { + err = c.Conn.SetReadDeadline(time.Now().Add(c.readDeadline)) + if err != nil { + return + } + + var messageType int + messageType, bytes, err = c.Conn.ReadMessage() + if messageType != websocket.BinaryMessage { + err = fmt.Errorf("only support websocket binary message") + return + } + return +} diff --git a/internal/gateway_ws/session/subject.go b/internal/gateway_ws/session/subject.go new file mode 100644 index 0000000..4f8ec2c --- /dev/null +++ b/internal/gateway_ws/session/subject.go @@ -0,0 +1,23 @@ +package session + +import "github.com/bytedance/sonic" + +// Subject 消息传递主体 +type Subject struct { + Uid string `json:"uid,omitempty"` + Online int8 `json:"online,omitempty"` // 0离线,1在线 + Time int64 `json:"time,omitempty"` + Gate string `json:"gate,omitempty"` // 连接的网关addr +} + +func (sub *Subject) MarshalBinary() (data []byte, err error) { + bytes, err := sonic.Marshal(sub) + if err != nil { + return nil, err + } + return bytes, nil +} + +func (sub *Subject) UnmarshalBinary(data []byte) error { + return sonic.Unmarshal(data, sub) +} diff --git a/internal/postal/logic/postal_server.go b/internal/postal/logic/postal_server.go index 697147d..0bedd4b 100644 --- a/internal/postal/logic/postal_server.go +++ b/internal/postal/logic/postal_server.go @@ -53,7 +53,7 @@ func (p *PostalServer) Run(conf config.GrpcConfig, postalRegister *discovery.Reg } // run serve - logger.Infof("%s server running %s\n", regConf.Name, listen.Addr().String()) + logger.Infof("%s grpc server running %s\n", regConf.Name, listen.Addr().String()) err = server.Serve(listen) return } diff --git a/internal/postal/ws/ws_server.go b/internal/postal/ws/ws_server.go deleted file mode 100644 index c76027c..0000000 --- a/internal/postal/ws/ws_server.go +++ /dev/null @@ -1,4 +0,0 @@ -package ws - -type WebsocketServer struct { -} diff --git a/pkg/grpc/discovery/resolver.go b/pkg/grpc/discovery/resolver.go index 2d59da7..e39b3fa 100644 --- a/pkg/grpc/discovery/resolver.go +++ b/pkg/grpc/discovery/resolver.go @@ -2,9 +2,9 @@ package discovery import ( "context" + "sonet/pkg/utils/logger" "time" - "github.com/sirupsen/logrus" clientv3 "go.etcd.io/etcd/client/v3" "google.golang.org/grpc/resolver" ) @@ -16,7 +16,6 @@ const ( // Resolver for grpc client type Resolver struct { schema string - EtcdAddrs []string DialTimeout int closeCh chan struct{} @@ -25,17 +24,15 @@ type Resolver struct { keyPrifix string srvAddrsList []resolver.Address - cc resolver.ClientConn - logger *logrus.Logger + cc resolver.ClientConn } // NewResolver create a new resolver.Builder base on etcd -func NewResolver(etcdAddrs []string, logger *logrus.Logger) *Resolver { +func NewResolver(client *clientv3.Client) *Resolver { return &Resolver{ + cli: client, schema: schema, - EtcdAddrs: etcdAddrs, DialTimeout: 3, - logger: logger, } } @@ -48,7 +45,7 @@ func (r *Resolver) Scheme() string { func (r *Resolver) Build(target resolver.Target, cc resolver.ClientConn, opts resolver.BuildOptions) (rr resolver.Resolver, err error) { r.cc = cc r.keyPrifix = BuildPrefix(Server{Name: target.Endpoint()}) - if _, err := r.start(); err != nil { + if err = r.start(); err != nil { return nil, err } return r, nil @@ -59,32 +56,30 @@ func (r *Resolver) ResolveNow(o resolver.ResolveNowOptions) {} // Close resolver.Resolver interface func (r *Resolver) Close() { + if r.closeCh == nil { + return + } r.closeCh <- struct{}{} + <-r.closeCh } // start -func (r *Resolver) start() (chan<- struct{}, error) { +func (r *Resolver) start() error { var err error - r.cli, err = clientv3.New(clientv3.Config{ - Endpoints: r.EtcdAddrs, - Username: "root", - Password: "sopod@etcd", - DialTimeout: time.Duration(r.DialTimeout) * time.Second, - }) - if err != nil { - return nil, err - } + resolver.Register(r) - r.closeCh = make(chan struct{}) + if r.closeCh == nil { + r.closeCh = make(chan struct{}) + } if err = r.sync(); err != nil { - return nil, err + return err } go r.watch() - return r.closeCh, nil + return nil } // watch update events @@ -95,6 +90,7 @@ func (r *Resolver) watch() { for { select { case <-r.closeCh: + r.closeCh <- struct{}{} return case res, ok := <-r.watchCh: if ok { @@ -102,7 +98,7 @@ func (r *Resolver) watch() { } case <-ticker.C: if err := r.sync(); err != nil { - r.logger.Error("sync failed", err) + logger.Error("resolver sync failed: ", err) } } } @@ -126,7 +122,10 @@ func (r *Resolver) update(events []*clientv3.Event) { } if !Exist(r.srvAddrsList, addr) { r.srvAddrsList = append(r.srvAddrsList, addr) - r.cc.UpdateState(resolver.State{Addresses: r.srvAddrsList}) + err = r.cc.UpdateState(resolver.State{Addresses: r.srvAddrsList}) + if err != nil { + logger.Error("resolver conn update put state err: ", err) + } } case clientv3.EventTypeDelete: info, err = SplitPath(string(ev.Kv.Key)) @@ -136,7 +135,10 @@ func (r *Resolver) update(events []*clientv3.Event) { addr := resolver.Address{Addr: info.Addr} if s, ok := Remove(r.srvAddrsList, addr); ok { r.srvAddrsList = s - r.cc.UpdateState(resolver.State{Addresses: r.srvAddrsList}) + err = r.cc.UpdateState(resolver.State{Addresses: r.srvAddrsList}) + if err != nil { + logger.Error("resolver conn update delete state err: ", err) + } } } } diff --git a/pkg/grpc/generic/generic_client.go b/pkg/grpc/generic/generic_client.go new file mode 100644 index 0000000..0bf2cf1 --- /dev/null +++ b/pkg/grpc/generic/generic_client.go @@ -0,0 +1,145 @@ +package generic + +import ( + "context" + "fmt" + "github.com/jhump/protoreflect/desc" + "github.com/jhump/protoreflect/dynamic" + "github.com/jhump/protoreflect/dynamic/grpcdynamic" + "github.com/jhump/protoreflect/grpcreflect" + "google.golang.org/grpc" + "google.golang.org/grpc/reflection/grpc_reflection_v1alpha" + "sonet/pkg/grpc/generic/desc_source" + "sonet/pkg/utils/logger" + "sync" +) + +type GrpcGenericClient struct { + serviceName string + conn *grpc.ClientConn + + descSource desc_source.DescriptorSource + serviceDesc *desc.ServiceDescriptor + callerCache *sync.Map +} + +func NewGpcGenericClient(serviceName string, conn *grpc.ClientConn) *GrpcGenericClient { + return &GrpcGenericClient{ + serviceName: serviceName, + conn: conn, + } +} + +func (c *GrpcGenericClient) Init(ctx context.Context) (err error) { + // fetch service description + refClient := grpcreflect.NewClientV1Alpha(ctx, grpc_reflection_v1alpha.NewServerReflectionClient(c.conn)) + desc_source.DescriptorSourceFromServer(ctx, refClient) + dsc, e1 := c.descSource.FindSymbol(c.serviceName) + if e1 != nil { + err = fmt.Errorf("service %s not found", c.serviceName) + if desc_source.IsNotFoundError(e1) { + logger.Errorf("target server not expose service %q in FindSymbol", c.serviceName) + return + } + logger.Errorf("failed to query for service descriptor %q: %v", c.serviceName, e1) + return + } + + sd, ok := dsc.(*desc.ServiceDescriptor) + if !ok { + err = fmt.Errorf("service %s not found", c.serviceName) + logger.Errorf("target server not expose service %q", c.serviceName) + return + } + c.serviceDesc = sd + c.callerCache = &sync.Map{} + return +} + +func (c *GrpcGenericClient) ServiceName() string { + return c.serviceName +} + +type methodCaller struct { + Mtd *desc.MethodDescriptor + MsgFactory *dynamic.MessageFactory + Stub grpcdynamic.Stub +} + +func (c *GrpcGenericClient) InvokeUnary(ctx context.Context, method string, reqBytes []byte, opts ...grpc.CallOption) (resp *dynamic.Message, err error) { + // cache method desc + caller, err := c.getMethodCaller(method) + if err != nil { + return + } + + reqMessage := caller.MsgFactory.NewMessage(caller.Mtd.GetInputType()) + if err = reqMessage.(*dynamic.Message).Unmarshal(reqBytes); err != nil { + err = fmt.Errorf("unmarshal req bytes error: %s", err.Error()) + return + } + + res, err := caller.Stub.InvokeRpc(ctx, caller.Mtd, reqMessage, opts...) + if err != nil { + return + } + resp = res.(*dynamic.Message) + return +} + +// getMethodCaller load generic resource from cache +func (c *GrpcGenericClient) getMethodCaller(method string) (caller *methodCaller, err error) { + val, ok := c.callerCache.Load(method) + if ok { + caller = val.(*methodCaller) + } else { + // method desc + caller = &methodCaller{} + caller.Mtd = c.serviceDesc.FindMethodByName(method) + if caller.Mtd == nil { + logger.Errorf("service %q does not include a method named %q", c.serviceName, method) + err = fmt.Errorf("method %s not found", method) + return + } + + // message factory + var ext dynamic.ExtensionRegistry + if err = c.fetchAllExtensions(&ext, caller.Mtd.GetInputType()); err != nil { + return + } + if err = c.fetchAllExtensions(&ext, caller.Mtd.GetOutputType()); err != nil { + return + } + caller.MsgFactory = dynamic.NewMessageFactoryWithExtensionRegistry(&ext) + + // stub + caller.Stub = grpcdynamic.NewStubWithMessageFactory(c.conn, caller.MsgFactory) + c.callerCache.Store(method, caller) + } + return +} + +func (c *GrpcGenericClient) fetchAllExtensions(ext *dynamic.ExtensionRegistry, md *desc.MessageDescriptor) (err error) { + msgTypeName := md.GetFullyQualifiedName() + if len(md.GetExtensionRanges()) > 0 { + fds, err := c.descSource.AllExtensionsForType(msgTypeName) + if err != nil { + return fmt.Errorf("failed to query for extensions of type %s: %v", msgTypeName, err) + } + for _, fd := range fds { + if err := ext.AddExtension(fd); err != nil { + return fmt.Errorf("could not register extension %s of type %s: %v", fd.GetFullyQualifiedName(), msgTypeName, err) + } + } + } + // recursively fetch extensions for the types of any message fields + for _, fd := range md.GetFields() { + if fd.GetMessageType() != nil { + err := c.fetchAllExtensions(ext, fd.GetMessageType()) + if err != nil { + return err + } + } + } + return nil +} diff --git a/pkg/grpc/generic/generic_client_factory.go b/pkg/grpc/generic/generic_client_factory.go new file mode 100644 index 0000000..6d48291 --- /dev/null +++ b/pkg/grpc/generic/generic_client_factory.go @@ -0,0 +1,51 @@ +package generic + +import ( + "context" + "fmt" + "google.golang.org/grpc" + "sonet/pkg/grpc/discovery" + "sync" +) + +type GrpcGenericClientFactory struct { + resolver *discovery.Resolver + defaultOpts []grpc.DialOption + clientCache *sync.Map +} + +func NewGpcGenericClientFactory(resolver *discovery.Resolver, defaultOpts ...grpc.DialOption) *GrpcGenericClientFactory { + return &GrpcGenericClientFactory{ + resolver: resolver, + defaultOpts: defaultOpts, + } +} + +func (f *GrpcGenericClientFactory) Init() { + f.clientCache = &sync.Map{} +} + +func (f *GrpcGenericClientFactory) NewClient(ctx context.Context, serviceName string, opts ...grpc.DialOption) (client *GrpcGenericClient, err error) { + addr := fmt.Sprintf("%s:///%s", f.resolver.Scheme(), serviceName) + dialOpts := make([]grpc.DialOption, len(f.defaultOpts)+len(opts)) + dialOpts = append(dialOpts, f.defaultOpts...) + dialOpts = append(dialOpts, opts...) + conn, err := grpc.DialContext(ctx, addr, dialOpts...) + client = NewGpcGenericClient(serviceName, conn) + err = client.Init(ctx) + return +} + +func (f *GrpcGenericClientFactory) GetClient(ctx context.Context, serviceName string, opts ...grpc.DialOption) (client *GrpcGenericClient, err error) { + val, ok := f.clientCache.Load(serviceName) + if ok { + client = val.(*GrpcGenericClient) + return + } + client, err = f.NewClient(ctx, serviceName, opts...) + if err != nil { + return + } + f.clientCache.Store(serviceName, client) + return +} diff --git a/pkg/protocol/protocol.go b/pkg/protocol/protocol.go index c450f77..596c326 100644 --- a/pkg/protocol/protocol.go +++ b/pkg/protocol/protocol.go @@ -1,38 +1,123 @@ package protocol +import ( + "encoding/binary" + "fmt" +) + // protocol // 1byte: magic: 99 -// 1byte: 1 request, 2 response, 3 event, 4 error +// 1byte: type 1 request, 2 response, 3 event, 4 error // 1byte: status: 20-OK, 30-CLIENT_TIMEOUT, 31-SERVER_TIMEOUT, 40-BAD_REQUEST, 41-BAD_RESPONSE, 44-SERVICE_NOT_FOUND, 50-CLIENT_ERROR, 51-SERVER_ERROR, 52-SERVICE_ERROR -// 4bit: svc,method url desc: 1 name, 2 number -// 4bit: serialize: 1proto, 2json +// 4bit: urlType svc,method url desc: 1 name, 2 number +// 4bit: serializeType: 1proto, 2json // 4byte: seqId // string: service, method / 4byte svc, 4byte method/4byte notice // proto bytes / json bytes const ( - Magic int8 = 99 + Magic byte = 99 - TypeRequest int8 = 1 - TypeResponse int8 = 2 - TypeNotice int8 = 3 - TypeError int8 = 4 + TypeRequest byte = 1 + TypeResponse byte = 2 + TypeNotice byte = 3 + TypeError byte = 4 ) type Header struct { - Magic int8 - Type int8 // 1 request, 2 response, 3 event, 4 error - Status int8 - UrlType int8 - SerializeType int8 + Magic byte + Type byte // 1 request, 2 response, 3 event, 4 error + Status byte + UrlType byte // 4bit: svc,method url desc: 1 name, 2 number + SerializeType byte // 4bit: body serialize: 1proto, 2json SeqId int32 - SvcNo int32 - MethodNo int32 - Svc string - Method string + SvcNo int32 // optional 1 + TargetNo int32 // optional 2 + Svc string // optional 1 + Target string // optional 2 method/message } type Payload struct { Header Header Body []byte } + +func Decode(bytes []byte) (payload *Payload, err error) { + header := Header{} + + header.Magic = bytes[0] + if header.Magic != Magic { + err = fmt.Errorf("unknown magic: %d", header.Magic) + return + } + header.Type = bytes[1] + header.Status = bytes[2] + header.UrlType = bytes[3] >> 4 + header.SerializeType = bytes[3] & 0xF + header.SeqId = int32(binary.BigEndian.Uint32(bytes[4:8])) + + var cursor int + switch header.UrlType { + case 1: + svcLen := int(binary.BigEndian.Uint32(bytes[8:12])) + header.Svc = string(bytes[12 : 12+svcLen]) + cursor = 12 + svcLen + + targetLen := int(binary.BigEndian.Uint32(bytes[cursor : cursor+4])) + header.Target = string(bytes[cursor+4 : cursor+4+targetLen]) + cursor = cursor + 4 + targetLen + case 2: + header.SvcNo = int32(binary.BigEndian.Uint32(bytes[8:12])) + header.TargetNo = int32(binary.BigEndian.Uint32(bytes[12:16])) + cursor = 16 + default: + err = fmt.Errorf("unknown url type: %d", header.UrlType) + return + } + + payload = &Payload{Header: header} + payload.Body = bytes[cursor:] + return +} + +func Encode(payload *Payload) (bytes []byte, err error) { + header := payload.Header + headerLen := 16 + var svc, target []byte + var svcLen, targetLen, bodyLen int + if header.UrlType == 1 { + svc = []byte(header.Svc) + target = []byte(header.Target) + svcLen = len(svc) + targetLen = len(target) + headerLen += svcLen + targetLen + } + bodyLen = len(payload.Body) + bytes = make([]byte, headerLen+bodyLen) + bytes[0] = header.Magic + bytes[1] = header.Type + bytes[2] = header.Status + bytes[3] = (header.UrlType << 4) | header.SerializeType + binary.BigEndian.PutUint32(bytes[4:8], uint32(header.SeqId)) + + var cursor int + switch header.UrlType { + case 1: + binary.BigEndian.PutUint32(bytes[8:12], uint32(svcLen)) + copy(bytes[12:12+svcLen], svc) + cursor = 12 + svcLen + binary.BigEndian.PutUint32(bytes[cursor:cursor+4], uint32(targetLen)) + copy(bytes[cursor+4:cursor+4+targetLen], target) + cursor = cursor + 4 + targetLen + case 2: + binary.BigEndian.PutUint32(bytes[8:12], uint32(header.SvcNo)) + binary.BigEndian.PutUint32(bytes[12:16], uint32(header.TargetNo)) + cursor = 16 + default: + err = fmt.Errorf("unknown url type: %d", header.UrlType) + } + if bodyLen > 0 { + copy(bytes[cursor:], payload.Body) + } + return +} diff --git a/pkg/protocol/protocol_test.go b/pkg/protocol/protocol_test.go new file mode 100644 index 0000000..f1193d7 --- /dev/null +++ b/pkg/protocol/protocol_test.go @@ -0,0 +1,50 @@ +package protocol + +import ( + "errors" + "testing" +) + +func TestProtocolCodec(t *testing.T) { + header := Header{ + Magic: Magic, + Type: 1, + Status: 20, + UrlType: 1, + SerializeType: 2, + SeqId: 10086, + Svc: "Postal", + Target: "Deliver", + } + payload := &Payload{ + Header: header, + Body: []byte(`{"receiver":"10001啊"}`), + } + bytes, err := Encode(payload) + if err != nil { + t.Error(err) + return + } + payload2, err := Decode(bytes) + if err != nil { + t.Error(err) + return + } + header2 := payload2.Header + body2Len := len(payload2.Body) + ok := header2.Magic == header.Magic && + header2.Status == header.Status && + header2.UrlType == header.UrlType && + header2.SerializeType == header.SerializeType && + header2.SeqId == header.SeqId && + header2.SvcNo == header.SvcNo && + header2.TargetNo == header.TargetNo && + header2.Svc == header.Svc && + header2.Target == header.Target && + body2Len == len(payload.Body) && + payload2.Body[0] == payload.Body[0] && + payload2.Body[body2Len-1] == payload.Body[body2Len-1] + if !ok { + t.Error(errors.New("decode not equals")) + } +} diff --git a/pkg/utils/resp/resp.go b/pkg/utils/resp/resp.go new file mode 100644 index 0000000..098e42a --- /dev/null +++ b/pkg/utils/resp/resp.go @@ -0,0 +1,61 @@ +package resp + +import ( + "github.com/bytedance/sonic" + "github.com/cloudwego/kitex/pkg/klog" +) + +const ( + CodeOK = 200 + CodeFail = 400 + CodeError = 500 +) + +type H map[string]interface{} + +// Response 响应体包装 +type Response struct { + Seq int `json:"seq,omitempty"` + Code int `json:"code,omitempty"` + Msg string `json:"msg,omitempty"` + Data any `json:"data,omitempty"` + Extra map[string]any `json:"extra,omitempty"` +} + +func (resp *Response) Json() []byte { + json, err := sonic.Marshal(resp) + if err != nil { + klog.Error("unknown json error: ", err) + return nil + } + return json +} + +func (resp *Response) JsonString() string { + return string(resp.Json()) +} + +func SeqResp(seq int, code int, msg string, data any) *Response { + return &Response{ + Seq: seq, + Code: code, + Msg: msg, + Data: data, + } +} + +func Resp(code int, msg string, data any) *Response { + return SeqResp(0, code, msg, data) +} + +func Success(data any) *Response { + return Resp(CodeOK, "", data) +} + +func Fail(msg string) *Response { + return Resp(CodeFail, msg, nil) +} + +func Error(msg string) *Response { + return Resp(CodeError, msg, nil) +}