master
Kevin Wan 4 years ago committed by GitHub
parent 38abfb80ed
commit 9602494454
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23

@ -1,10 +1,22 @@
package internal package internal
import "github.com/tal-tech/go-zero/core/discov" import (
"os"
"strings"
"github.com/tal-tech/go-zero/core/discov"
"github.com/tal-tech/go-zero/core/netx"
)
const (
allEths = "0.0.0.0"
envPodIp = "POD_IP"
)
func NewRpcPubServer(etcdEndpoints []string, etcdKey, listenOn string, opts ...ServerOption) (Server, error) { func NewRpcPubServer(etcdEndpoints []string, etcdKey, listenOn string, opts ...ServerOption) (Server, error) {
registerEtcd := func() error { registerEtcd := func() error {
pubClient := discov.NewPublisher(etcdEndpoints, etcdKey, listenOn) pubListenOn := figureOutListenOn(listenOn)
pubClient := discov.NewPublisher(etcdEndpoints, etcdKey, pubListenOn)
return pubClient.KeepAlive() return pubClient.KeepAlive()
} }
server := keepAliveServer{ server := keepAliveServer{
@ -27,3 +39,25 @@ func (ags keepAliveServer) Start(fn RegisterFn) error {
return ags.Server.Start(fn) return ags.Server.Start(fn)
} }
func figureOutListenOn(listenOn string) string {
fields := strings.Split(listenOn, ":")
if len(fields) == 0 {
return listenOn
}
host := fields[0]
if len(host) > 0 && host != allEths {
return listenOn
}
ip := os.Getenv(envPodIp)
if len(ip) == 0 {
ip = netx.InternalIp()
}
if len(ip) == 0 {
return listenOn
}
return strings.Join(append([]string{ip}, fields[1:]...), ":")
}

@ -0,0 +1,33 @@
package internal
import (
"testing"
"github.com/stretchr/testify/assert"
"github.com/tal-tech/go-zero/core/netx"
)
func TestFigureOutListenOn(t *testing.T) {
tests := []struct {
input string
expect string
}{
{
input: "192.168.0.5:1234",
expect: "192.168.0.5:1234",
},
{
input: "0.0.0.0:8080",
expect: netx.InternalIp() + ":8080",
},
{
input: ":8080",
expect: netx.InternalIp() + ":8080",
},
}
for _, test := range tests {
val := figureOutListenOn(test.input)
assert.Equal(t, test.expect, val)
}
}

@ -2,13 +2,10 @@ package zrpc
import ( import (
"log" "log"
"os"
"strings"
"time" "time"
"github.com/tal-tech/go-zero/core/load" "github.com/tal-tech/go-zero/core/load"
"github.com/tal-tech/go-zero/core/logx" "github.com/tal-tech/go-zero/core/logx"
"github.com/tal-tech/go-zero/core/netx"
"github.com/tal-tech/go-zero/core/stat" "github.com/tal-tech/go-zero/core/stat"
"github.com/tal-tech/go-zero/zrpc/internal" "github.com/tal-tech/go-zero/zrpc/internal"
"github.com/tal-tech/go-zero/zrpc/internal/auth" "github.com/tal-tech/go-zero/zrpc/internal/auth"
@ -16,11 +13,6 @@ import (
"google.golang.org/grpc" "google.golang.org/grpc"
) )
const (
allEths = "0.0.0.0"
envPodIp = "POD_IP"
)
type RpcServer struct { type RpcServer struct {
server internal.Server server internal.Server
register internal.RegisterFn register internal.RegisterFn
@ -44,8 +36,7 @@ func NewServer(c RpcServerConf, register internal.RegisterFn) (*RpcServer, error
var server internal.Server var server internal.Server
metrics := stat.NewMetrics(c.ListenOn) metrics := stat.NewMetrics(c.ListenOn)
if c.HasEtcd() { if c.HasEtcd() {
listenOn := figureOutListenOn(c.ListenOn) server, err = internal.NewRpcPubServer(c.Etcd.Hosts, c.Etcd.Key, c.ListenOn, internal.WithMetrics(metrics))
server, err = internal.NewRpcPubServer(c.Etcd.Hosts, c.Etcd.Key, listenOn, internal.WithMetrics(metrics))
if err != nil { if err != nil {
return nil, err return nil, err
} }
@ -92,28 +83,6 @@ func (rs *RpcServer) Stop() {
logx.Close() logx.Close()
} }
func figureOutListenOn(listenOn string) string {
fields := strings.Split(listenOn, ":")
if len(fields) == 0 {
return listenOn
}
host := fields[0]
if len(host) > 0 && host != allEths {
return listenOn
}
ip := os.Getenv(envPodIp)
if len(ip) == 0 {
ip = netx.InternalIp()
}
if len(ip) == 0 {
return listenOn
}
return strings.Join(append([]string{ip}, fields[1:]...), ":")
}
func setupInterceptors(server internal.Server, c RpcServerConf, metrics *stat.Metrics) error { func setupInterceptors(server internal.Server, c RpcServerConf, metrics *stat.Metrics) error {
if c.CpuThreshold > 0 { if c.CpuThreshold > 0 {
shedder := load.NewAdaptiveShedder(load.WithCpuThreshold(c.CpuThreshold)) shedder := load.NewAdaptiveShedder(load.WithCpuThreshold(c.CpuThreshold))

@ -7,7 +7,6 @@ import (
"github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert"
"github.com/tal-tech/go-zero/core/discov" "github.com/tal-tech/go-zero/core/discov"
"github.com/tal-tech/go-zero/core/logx" "github.com/tal-tech/go-zero/core/logx"
"github.com/tal-tech/go-zero/core/netx"
"github.com/tal-tech/go-zero/core/service" "github.com/tal-tech/go-zero/core/service"
"github.com/tal-tech/go-zero/core/stat" "github.com/tal-tech/go-zero/core/stat"
"github.com/tal-tech/go-zero/core/stores/redis" "github.com/tal-tech/go-zero/core/stores/redis"
@ -16,31 +15,6 @@ import (
"google.golang.org/grpc" "google.golang.org/grpc"
) )
func TestFigureOutListenOn(t *testing.T) {
tests := []struct {
input string
expect string
}{
{
input: "192.168.0.5:1234",
expect: "192.168.0.5:1234",
},
{
input: "0.0.0.0:8080",
expect: netx.InternalIp() + ":8080",
},
{
input: ":8080",
expect: netx.InternalIp() + ":8080",
},
}
for _, test := range tests {
val := figureOutListenOn(test.input)
assert.Equal(t, test.expect, val)
}
}
func TestServer_setupInterceptors(t *testing.T) { func TestServer_setupInterceptors(t *testing.T) {
server := new(mockedServer) server := new(mockedServer)
err := setupInterceptors(server, RpcServerConf{ err := setupInterceptors(server, RpcServerConf{

Loading…
Cancel
Save