This is an automated email from the ASF dual-hosted git repository.
Alanxtl pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/dubbo-go.git
The following commit(s) were added to refs/heads/develop by this push:
new 39eeb2bda fix(protocol, registry, remoting, graceful_shutdown): handle
ignored errors across rpc paths (#3574)
39eeb2bda is described below
commit 39eeb2bda23613fad65545dfb05de9b45209b5e7
Author: Nene7ko_ <[email protected]>
AuthorDate: Tue Aug 11 18:51:35 2026 +0800
fix(protocol, registry, remoting, graceful_shutdown): handle ignored errors
across rpc paths (#3574)
* fix: handle ignored errors across rpc paths
* style: fix getty import formatting
* fix: avoid shadowing hessian encode error
* fix: address PR 3574 review feedback
* fix: satisfy testifylint error assertion
---
graceful_shutdown/shutdown.go | 3 +-
protocol/dubbo/hessian2/hessian_request.go | 66 +++++-----
protocol/dubbo/hessian2/hessian_request_test.go | 48 ++++++++
protocol/dubbo/impl/hessian.go | 152 +++++++++++++++---------
protocol/dubbo/impl/hessian_test.go | 32 +++++
protocol/grpc/grpc_invoker.go | 6 +-
protocol/grpc/grpc_invoker_test.go | 28 +++++
protocol/triple/server.go | 28 ++++-
registry/nacos/listener.go | 4 +-
remoting/getty/config.go | 10 +-
remoting/polaris/parser/parser.go | 5 +-
11 files changed, 284 insertions(+), 98 deletions(-)
diff --git a/graceful_shutdown/shutdown.go b/graceful_shutdown/shutdown.go
index f01d16db3..b024a1f29 100644
--- a/graceful_shutdown/shutdown.go
+++ b/graceful_shutdown/shutdown.go
@@ -422,7 +422,8 @@ func invokeCustomShutdownCallback(timeout time.Duration,
callback func()) {
func getProtocolSafely(name string) (protocol protocolbase.Protocol, ok bool) {
defer func() {
- if recover() != nil {
+ if recovered := recover(); recovered != nil {
+ logger.Warnf("[GracefulShutdown] get protocol %s
panicked, err=%v", name, recovered)
protocol = nil
ok = false
}
diff --git a/protocol/dubbo/hessian2/hessian_request.go
b/protocol/dubbo/hessian2/hessian_request.go
index b503295fa..295ff10c2 100644
--- a/protocol/dubbo/hessian2/hessian_request.go
+++ b/protocol/dubbo/hessian2/hessian_request.go
@@ -85,8 +85,6 @@ func EnsureRequest(body any) *DubboRequest {
func packRequest(service Service, header DubboHeader, req any) ([]byte, error)
{
var (
- err error
- types string
byteArray []byte
pkgLen int
)
@@ -126,31 +124,53 @@ func packRequest(service Service, header DubboHeader, req
any) ([]byte, error) {
// body
//////////////////////////////////////////
if hb {
- _ = encoder.Encode(nil)
- goto END
+ if err := encoder.Encode(nil); err != nil {
+ return nil, perrors.Wrap(err, "failed to encode
heartbeat request")
+ }
+ } else {
+ if err := encodeRequestBody(encoder, service, request, args);
err != nil {
+ return nil, err
+ }
}
+ byteArray = encoder.Buffer()
+ pkgLen = len(byteArray)
+ if pkgLen > int(DEFAULT_LEN) { // recommand 8M
+ logger.Warnf("[Dubbo][Hessian2] data length %d too large,
recommand max payload %d. "+
+ "Dubbo java can't handle the package whose size is
greater than %d!!!", pkgLen, DEFAULT_LEN, DEFAULT_LEN)
+ }
+ // byteArray{body length}
+ binary.BigEndian.PutUint32(byteArray[12:], uint32(pkgLen-HEADER_LENGTH))
+ return byteArray, nil
+}
+
+func encodeRequestBody(encoder *hessian.Encoder, service Service, request
*DubboRequest, args []any) error {
// dubbo version + path + version + method
- if err = encoder.Encode(DEFAULT_DUBBO_PROTOCOL_VERSION); err != nil {
- logger.Warnf("[Dubbo][Hessian2] encode default dubbo protocol
version failed, err=%v", err)
+ if err := encoder.Encode(DEFAULT_DUBBO_PROTOCOL_VERSION); err != nil {
+ return perrors.Wrap(err, "failed to encode default dubbo
protocol version")
}
- if err = encoder.Encode(service.Path); err != nil {
- logger.Warnf("[Dubbo][Hessian2] encode service path failed,
err=%v", err)
+ if err := encoder.Encode(service.Path); err != nil {
+ return perrors.Wrap(err, "failed to encode service path")
}
- if err = encoder.Encode(service.Version); err != nil {
- logger.Warnf("[Dubbo][Hessian2] encode service version failed,
err=%v", err)
+ if err := encoder.Encode(service.Version); err != nil {
+ return perrors.Wrap(err, "failed to encode service version")
}
- if err = encoder.Encode(service.Method); err != nil {
- logger.Warnf("[Dubbo][Hessian2] encode service method failed,
err=%v", err)
+ if err := encoder.Encode(service.Method); err != nil {
+ return perrors.Wrap(err, "failed to encode service method")
}
// args = args type list + args value list
- if types, err = getArgsTypeList(args); err != nil {
- return nil, perrors.Wrapf(err, " PackRequest(args:%+v)", args)
+ types, err := getArgsTypeList(args)
+ if err != nil {
+ return perrors.Wrapf(err, " PackRequest(args:%+v)", args)
+ }
+ if err := encoder.Encode(types); err != nil {
+ return perrors.Wrap(err, "failed to encode argument types")
}
- _ = encoder.Encode(types)
for _, v := range args {
- _ = encoder.Encode(v)
+ if err := encoder.Encode(v); err != nil {
+ return perrors.Wrapf(err, "failed to encode argument of
type %T", v)
+ }
}
request.Attachments[PATH_KEY] = service.Path
@@ -165,18 +185,10 @@ func packRequest(service Service, header DubboHeader, req
any) ([]byte, error) {
request.Attachments[TIMEOUT_KEY] =
strconv.Itoa(int(service.Timeout / time.Millisecond))
}
- _ = encoder.Encode(request.Attachments)
-
-END:
- byteArray = encoder.Buffer()
- pkgLen = len(byteArray)
- if pkgLen > int(DEFAULT_LEN) { // recommand 8M
- logger.Warnf("[Dubbo][Hessian2] data length %d too large,
recommand max payload %d. "+
- "Dubbo java can't handle the package whose size is
greater than %d!!!", pkgLen, DEFAULT_LEN, DEFAULT_LEN)
+ if err := encoder.Encode(request.Attachments); err != nil {
+ return perrors.Wrap(err, "failed to encode request attachments")
}
- // byteArray{body length}
- binary.BigEndian.PutUint32(byteArray[12:], uint32(pkgLen-HEADER_LENGTH))
- return byteArray, nil
+ return nil
}
// hessian decode request body
diff --git a/protocol/dubbo/hessian2/hessian_request_test.go
b/protocol/dubbo/hessian2/hessian_request_test.go
index 7a9d70ead..5d55eeb91 100644
--- a/protocol/dubbo/hessian2/hessian_request_test.go
+++ b/protocol/dubbo/hessian2/hessian_request_test.go
@@ -18,6 +18,7 @@
package hessian2
import (
+ "encoding/binary"
"reflect"
"strconv"
"testing"
@@ -90,6 +91,53 @@ func TestPackRequest(t *testing.T) {
}
}
+func TestPackRequestHeartbeat(t *testing.T) {
+ data, err := packRequest(Service{}, DubboHeader{Type:
PackageHeartbeat}, NewRequest([]any{}, nil))
+
+ require.NoError(t, err)
+ require.Len(t, data, HEADER_LENGTH+1)
+ assert.Equal(t, DubboRequestHeartbeatHeader[:HEADER_LENGTH-4],
data[:HEADER_LENGTH-4])
+ assert.Equal(t, uint32(1),
binary.BigEndian.Uint32(data[HEADER_LENGTH-4:HEADER_LENGTH]))
+ assert.Equal(t, byte('N'), data[HEADER_LENGTH])
+}
+
+func TestPackRequestReturnsEncodeErrors(t *testing.T) {
+ tests := []struct {
+ name string
+ request *DubboRequest
+ errString string
+ }{
+ {
+ name: "unsupported map key argument",
+ request: NewRequest([]any{
+ map[complex64]string{1 + 2i: "value"},
+ }, nil),
+ errString: "failed to encode argument of type
map[complex64]string",
+ },
+ {
+ name: "unsupported attachment",
+ request: NewRequest([]any{}, map[string]any{
+ "unsupported": func() {},
+ }),
+ errString: "failed to encode request attachments",
+ },
+ }
+
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ data, err := packRequest(Service{
+ Path: "test",
+ Version: "v1.0",
+ Method: "test",
+ }, DubboHeader{Type: PackageRequest}, tt.request)
+
+ require.Error(t, err)
+ require.ErrorContains(t, err, tt.errString)
+ assert.Nil(t, data)
+ })
+ }
+}
+
func TestGetArgsTypeList(t *testing.T) {
type Test struct{}
str, err := getArgsTypeList([]any{nil, 1, []int{2}, true,
[]bool{false}, "a", []string{"b"}, Test{}, &Test{}, []Test{},
map[string]Test{}, TestEnumGender(MAN)})
diff --git a/protocol/dubbo/impl/hessian.go b/protocol/dubbo/impl/hessian.go
index becdfcdee..4a0e4d6e2 100644
--- a/protocol/dubbo/impl/hessian.go
+++ b/protocol/dubbo/impl/hessian.go
@@ -61,60 +61,9 @@ func (h HessianSerializer) Unmarshal(input []byte, p
*DubboPackage) error {
}
func marshalResponse(encoder *hessian.Encoder, p DubboPackage) ([]byte, error)
{
- header := p.Header
response := EnsureResponsePayload(p.Body)
- if header.ResponseStatus == Response_OK {
- if p.IsHeartBeat() {
- _ = encoder.Encode(nil)
- } else {
- var version string
- if attachmentVersion, ok :=
response.Attachments[DUBBO_VERSION_KEY]; ok {
- version = attachmentVersion.(string)
- }
- atta := isSupportResponseAttachment(version)
-
- var resWithException, resValue, resNullValue int32
- if atta {
- resWithException =
RESPONSE_WITH_EXCEPTION_WITH_ATTACHMENTS
- resValue = RESPONSE_VALUE_WITH_ATTACHMENTS
- resNullValue =
RESPONSE_NULL_VALUE_WITH_ATTACHMENTS
- } else {
- resWithException = RESPONSE_WITH_EXCEPTION
- resValue = RESPONSE_VALUE
- resNullValue = RESPONSE_NULL_VALUE
- }
-
- if response.Exception != nil { // throw error
- _ = encoder.Encode(resWithException)
- switch ex := response.Exception.(type) {
- case *hessian.GenericException:
- _ =
encoder.Encode(java_exception.NewDubboGenericException(ex.ExceptionClass,
ex.ExceptionMessage))
- case hessian.GenericException:
- _ =
encoder.Encode(java_exception.NewDubboGenericException(ex.ExceptionClass,
ex.ExceptionMessage))
- case java_exception.Throwabler:
- _ = encoder.Encode(ex)
- default:
- _ =
encoder.Encode(java_exception.NewThrowable(response.Exception.Error()))
- }
- } else {
- if response.RspObj == nil {
- _ = encoder.Encode(resNullValue)
- } else {
- _ = encoder.Encode(resValue)
- _ = encoder.Encode(response.RspObj) //
result
- }
- }
-
- if atta {
- _ = encoder.Encode(response.Attachments) //
attachments
- }
- }
- } else {
- if response.Exception != nil { // throw error
- _ = encoder.Encode(response.Exception.Error())
- } else {
- _ = encoder.Encode(response.RspObj)
- }
+ if err := encodeResponse(encoder, p, response); err != nil {
+ return nil, err
}
bs := encoder.Buffer()
// encNull
@@ -122,13 +71,97 @@ func marshalResponse(encoder *hessian.Encoder, p
DubboPackage) ([]byte, error) {
return bs, nil
}
+func encodeResponse(encoder *hessian.Encoder, p DubboPackage, response
*ResponsePayload) error {
+ if p.Header.ResponseStatus != Response_OK {
+ if response.Exception != nil {
+ return encodeHessianValue(encoder,
response.Exception.Error(), "response exception message")
+ }
+ return encodeHessianValue(encoder, response.RspObj, "response
value")
+ }
+ if p.IsHeartBeat() {
+ return encodeHessianValue(encoder, nil, "heartbeat response")
+ }
+
+ var version string
+ if attachmentVersion, ok := response.Attachments[DUBBO_VERSION_KEY]; ok
{
+ version = attachmentVersion.(string)
+ }
+ withAttachments := isSupportResponseAttachment(version)
+ resWithException, resValue, resNullValue :=
responseTypes(withAttachments)
+
+ if err := encodeResponsePayload(encoder, response, resWithException,
resValue, resNullValue); err != nil {
+ return err
+ }
+ if withAttachments {
+ return encodeHessianValue(encoder, response.Attachments,
"response attachments")
+ }
+ return nil
+}
+
+func encodeResponsePayload(encoder *hessian.Encoder, response
*ResponsePayload, resWithException, resValue, resNullValue int32) error {
+ if response.Exception != nil {
+ return encodeExceptionResponse(encoder, response.Exception,
resWithException)
+ }
+ if response.RspObj == nil {
+ return encodeHessianValue(encoder, resNullValue, "null response
type")
+ }
+ if err := encodeHessianValue(encoder, resValue, "response value type");
err != nil {
+ return err
+ }
+ return encodeHessianValue(encoder, response.RspObj, "response value")
+}
+
+func responseTypes(withAttachments bool) (int32, int32, int32) {
+ if withAttachments {
+ return RESPONSE_WITH_EXCEPTION_WITH_ATTACHMENTS,
RESPONSE_VALUE_WITH_ATTACHMENTS, RESPONSE_NULL_VALUE_WITH_ATTACHMENTS
+ }
+ return RESPONSE_WITH_EXCEPTION, RESPONSE_VALUE, RESPONSE_NULL_VALUE
+}
+
+func encodeResponseException(encoder *hessian.Encoder, exception error) error {
+ var value any
+ switch ex := exception.(type) {
+ case *hessian.GenericException:
+ value =
java_exception.NewDubboGenericException(ex.ExceptionClass, ex.ExceptionMessage)
+ case hessian.GenericException:
+ value =
java_exception.NewDubboGenericException(ex.ExceptionClass, ex.ExceptionMessage)
+ case java_exception.Throwabler:
+ value = ex
+ default:
+ value = java_exception.NewThrowable(exception.Error())
+ }
+ return encodeHessianValue(encoder, value, "response exception")
+}
+
+func encodeExceptionResponse(encoder *hessian.Encoder, exception error,
responseType int32) error {
+ if err := encodeHessianValue(encoder, responseType, "response exception
type"); err != nil {
+ return err
+ }
+ return encodeResponseException(encoder, exception)
+}
+
+func encodeHessianValue(encoder *hessian.Encoder, value any, description
string) error {
+ if err := encoder.Encode(value); err != nil {
+ return perrors.Wrapf(err, "failed to encode %s", description)
+ }
+ return nil
+}
+
func marshalRequest(encoder *hessian.Encoder, p DubboPackage) ([]byte, error) {
service := p.Service
request := EnsureRequestPayload(p.Body)
- _ = encoder.Encode(DEFAULT_DUBBO_PROTOCOL_VERSION)
- _ = encoder.Encode(service.Path)
- _ = encoder.Encode(service.Version)
- _ = encoder.Encode(service.Method)
+ if err := encodeHessianValue(encoder, DEFAULT_DUBBO_PROTOCOL_VERSION,
"dubbo protocol version"); err != nil {
+ return nil, err
+ }
+ if err := encodeHessianValue(encoder, service.Path, "service path");
err != nil {
+ return nil, err
+ }
+ if err := encodeHessianValue(encoder, service.Version, "service
version"); err != nil {
+ return nil, err
+ }
+ if err := encodeHessianValue(encoder, service.Method, "service
method"); err != nil {
+ return nil, err
+ }
args, ok := request.Params.([]any)
@@ -140,7 +173,10 @@ func marshalRequest(encoder *hessian.Encoder, p
DubboPackage) ([]byte, error) {
if err != nil {
return nil, perrors.Wrapf(err, " PackRequest(args:%+v)", args)
}
- _ = encoder.Encode(types)
+ err = encodeHessianValue(encoder, types, "argument types")
+ if err != nil {
+ return nil, err
+ }
for _, v := range args {
if e := encoder.Encode(v); e != nil {
return nil, perrors.Wrapf(e, "failed to encode
argument: %v", v)
diff --git a/protocol/dubbo/impl/hessian_test.go
b/protocol/dubbo/impl/hessian_test.go
index 16e29b424..51bb412c2 100644
--- a/protocol/dubbo/impl/hessian_test.go
+++ b/protocol/dubbo/impl/hessian_test.go
@@ -449,6 +449,38 @@ func TestMarshalResponse(t *testing.T) {
assert.NotNil(t, data)
})
+ t.Run("response with unsupported value", func(t *testing.T) {
+ encoder := hessian.NewEncoder()
+ pkg := DubboPackage{
+ Header: DubboHeader{ResponseStatus: Response_OK},
+ Body: &ResponsePayload{
+ RspObj: func() {},
+ Attachments: map[string]any{},
+ },
+ }
+
+ data, err := marshalResponse(encoder, pkg)
+ require.Error(t, err)
+ assert.Nil(t, data)
+ })
+
+ t.Run("response with unsupported attachment", func(t *testing.T) {
+ encoder := hessian.NewEncoder()
+ pkg := DubboPackage{
+ Header: DubboHeader{ResponseStatus: Response_OK},
+ Body: &ResponsePayload{
+ Attachments: map[string]any{
+ DUBBO_VERSION_KEY: "2.7.0",
+ "unsupported": func() {},
+ },
+ },
+ }
+
+ data, err := marshalResponse(encoder, pkg)
+ require.Error(t, err)
+ assert.Nil(t, data)
+ })
+
t.Run("response with value", func(t *testing.T) {
encoder := hessian.NewEncoder()
pkg := DubboPackage{
diff --git a/protocol/grpc/grpc_invoker.go b/protocol/grpc/grpc_invoker.go
index 9ab74ce68..1df29000e 100644
--- a/protocol/grpc/grpc_invoker.go
+++ b/protocol/grpc/grpc_invoker.go
@@ -123,8 +123,10 @@ func (gi *GrpcInvoker) Invoke(ctx context.Context,
invocation base.Invocation) r
// check err
if !res[1].IsNil() {
result.SetError(res[1].Interface().(error))
- } else {
- _ = hessian2.ReflectResponse(res[0], invocation.Reply())
+ } else if invocation.Reply() != nil {
+ if err := hessian2.ReflectResponse(res[0], invocation.Reply());
err != nil {
+ result.SetError(errors.WithStack(err))
+ }
}
return &result
diff --git a/protocol/grpc/grpc_invoker_test.go
b/protocol/grpc/grpc_invoker_test.go
index b1cb637fa..f59903ee3 100644
--- a/protocol/grpc/grpc_invoker_test.go
+++ b/protocol/grpc/grpc_invoker_test.go
@@ -20,6 +20,7 @@ package grpc
import (
"context"
"reflect"
+ "sync"
"testing"
)
@@ -31,11 +32,38 @@ import (
import (
"dubbo.apache.org/dubbo-go/v3/common"
"dubbo.apache.org/dubbo-go/v3/common/constant"
+ "dubbo.apache.org/dubbo-go/v3/protocol/base"
"dubbo.apache.org/dubbo-go/v3/protocol/grpc/internal/helloworld"
"dubbo.apache.org/dubbo-go/v3/protocol/grpc/internal/routeguide"
"dubbo.apache.org/dubbo-go/v3/protocol/invocation"
)
+type grpcInvokerTestClient struct{}
+
+func (*grpcInvokerTestClient) Call(context.Context) (string, error) {
+ return "response", nil
+}
+
+func TestUnaryInvokeReturnsReflectResponseError(t *testing.T) {
+ url := common.NewURLWithOptions(common.WithProtocol("grpc"))
+ invoker := &GrpcInvoker{
+ BaseInvoker: *base.NewBaseInvoker(url),
+ clientGuard: &sync.RWMutex{},
+ client: &Client{
+ invoker: reflect.ValueOf(&grpcInvokerTestClient{}),
+ },
+ }
+ invo := invocation.NewRPCInvocationWithOptions(
+ invocation.WithMethodName("Call"),
+ invocation.WithReply("reply must be a pointer"),
+ )
+
+ res := invoker.Invoke(context.Background(), invo)
+
+ require.Error(t, res.Error())
+ assert.ErrorContains(t, res.Error(), "@out should be a pointer")
+}
+
const (
helloworldURL =
"grpc://127.0.0.1:30000/GrpcGreeterImpl?accesslog=&anyhost=true&app.version=0.0.1&application=BDTService&async=false&bean.name=GrpcGreeterImpl"
+
"&category=providers&cluster=failover&dubbo=dubbo-provider-golang-2.6.0&environment=dev&execute.limit=&execute.limit.rejected.handler=&generic=false&group=&interface=io.grpc.examples.helloworld.GreeterGrpc%24IGreeter"
+
diff --git a/protocol/triple/server.go b/protocol/triple/server.go
index e6f3d894d..b1963cccf 100644
--- a/protocol/triple/server.go
+++ b/protocol/triple/server.go
@@ -492,7 +492,9 @@ func (s *Server) compatRegisterHandler(interfaceName
string, svc dubbo3.Dubbo3Gr
// Please refer to
protocol/triple/internal/proto/triple_gen/greettriple for procedure examples
// Error could be ignored because base is empty string
procedure := joinProcedure(interfaceName, method.MethodName)
- _ = s.triServer.RegisterCompatUnaryHandler(procedure,
method.MethodName, svc, tri.MethodHandler(method.Handler), opts...)
+ if err := s.triServer.RegisterCompatUnaryHandler(procedure,
method.MethodName, svc, tri.MethodHandler(method.Handler), opts...); err != nil
{
+ logger.Errorf("[Triple][Server] register compat unary
handler failed, procedure=%s, err=%v", procedure, err)
+ }
}
// Init stream handlers
@@ -509,7 +511,9 @@ func (s *Server) compatRegisterHandler(interfaceName
string, svc dubbo3.Dubbo3Gr
case stream.ServerStreams:
typ = tri.StreamTypeServer
}
- _ = s.triServer.RegisterCompatStreamHandler(procedure, svc,
typ, stream.Handler, opts...)
+ if err := s.triServer.RegisterCompatStreamHandler(procedure,
svc, typ, stream.Handler, opts...); err != nil {
+ logger.Errorf("[Triple][Server] register compat stream
handler failed, procedure=%s, err=%v", procedure, err)
+ }
}
}
@@ -539,7 +543,7 @@ func (s *Server) registerMethodHandler(procedure string, m
common.MethodInfo, in
}
func (s *Server) registerUnaryMethodHandler(procedure string, m
common.MethodInfo, invoker base.Invoker, opts ...tri.HandlerOption) {
- _ = s.triServer.RegisterUnaryHandler(
+ err := s.triServer.RegisterUnaryHandler(
procedure,
m.ReqInitFunc,
func(ctx context.Context, req *tri.Request) (*tri.Response,
error) {
@@ -557,10 +561,13 @@ func (s *Server) registerUnaryMethodHandler(procedure
string, m common.MethodInf
},
opts...,
)
+ if err != nil {
+ logger.Errorf("[Triple][Server] register unary handler failed,
procedure=%s, err=%v", procedure, err)
+ }
}
func (s *Server) registerClientStreamMethodHandler(procedure string, m
common.MethodInfo, invoker base.Invoker, opts ...tri.HandlerOption) {
- _ = s.triServer.RegisterClientStreamHandler(
+ err := s.triServer.RegisterClientStreamHandler(
procedure,
func(ctx context.Context, stream *tri.ClientStream)
(*tri.Response, error) {
args := []any{m.StreamInitFunc(stream)}
@@ -573,10 +580,13 @@ func (s *Server)
registerClientStreamMethodHandler(procedure string, m common.Me
},
opts...,
)
+ if err != nil {
+ logger.Errorf("[Triple][Server] register client stream handler
failed, procedure=%s, err=%v", procedure, err)
+ }
}
func (s *Server) registerServerStreamMethodHandler(procedure string, m
common.MethodInfo, invoker base.Invoker, opts ...tri.HandlerOption) {
- _ = s.triServer.RegisterServerStreamHandler(
+ err := s.triServer.RegisterServerStreamHandler(
procedure,
m.ReqInitFunc,
func(ctx context.Context, req *tri.Request, stream
*tri.ServerStream) error {
@@ -590,10 +600,13 @@ func (s *Server)
registerServerStreamMethodHandler(procedure string, m common.Me
},
opts...,
)
+ if err != nil {
+ logger.Errorf("[Triple][Server] register server stream handler
failed, procedure=%s, err=%v", procedure, err)
+ }
}
func (s *Server) registerBidiStreamMethodHandler(procedure string, m
common.MethodInfo, invoker base.Invoker, opts ...tri.HandlerOption) {
- _ = s.triServer.RegisterBidiStreamHandler(
+ err := s.triServer.RegisterBidiStreamHandler(
procedure,
func(ctx context.Context, stream *tri.BidiStream) error {
args := []any{m.StreamInitFunc(stream)}
@@ -606,6 +619,9 @@ func (s *Server) registerBidiStreamMethodHandler(procedure
string, m common.Meth
},
opts...,
)
+ if err != nil {
+ logger.Errorf("[Triple][Server] register bidi stream handler
failed, procedure=%s, err=%v", procedure, err)
+ }
}
func extractUnaryInvocationArgs(msg any) []any {
diff --git a/registry/nacos/listener.go b/registry/nacos/listener.go
index 8a5971642..99687b0e6 100644
--- a/registry/nacos/listener.go
+++ b/registry/nacos/listener.go
@@ -239,7 +239,9 @@ func (nl *nacosListener) Next() (*registry.ServiceEvent,
error) {
// Close stops the subscription and releases resources.
func (nl *nacosListener) Close() {
nl.once.Do(func() {
- _ = nl.stopListen()
+ if err := nl.stopListen(); err != nil {
+ logger.Warnf("[Registry][Nacos] unsubscribe listener
failed, serviceName=%s, err=%v", nl.serviceName, err)
+ }
close(nl.done)
})
}
diff --git a/remoting/getty/config.go b/remoting/getty/config.go
index 5143e075f..ec2322584 100644
--- a/remoting/getty/config.go
+++ b/remoting/getty/config.go
@@ -24,6 +24,8 @@ import (
import (
getty "github.com/apache/dubbo-getty"
+ "github.com/dubbogo/gost/log/logger"
+
perrors "github.com/pkg/errors"
)
@@ -140,7 +142,9 @@ func GetDefaultClientConfig() *ClientConfig {
SessionName: "client",
},
}
- _ = defaultClientConfig.CheckValidity()
+ if err := defaultClientConfig.CheckValidity(); err != nil {
+ logger.Errorf("[Remoting][Getty] invalid default client config,
err=%v", err)
+ }
return defaultClientConfig
}
@@ -166,7 +170,9 @@ func GetDefaultServerConfig() *ServerConfig {
SessionName: "server",
},
}
- _ = defaultServerConfig.CheckValidity()
+ if err := defaultServerConfig.CheckValidity(); err != nil {
+ logger.Errorf("[Remoting][Getty] invalid default server config,
err=%v", err)
+ }
return defaultServerConfig
}
diff --git a/remoting/polaris/parser/parser.go
b/remoting/polaris/parser/parser.go
index 7271f7478..835f377bc 100644
--- a/remoting/polaris/parser/parser.go
+++ b/remoting/polaris/parser/parser.go
@@ -88,7 +88,10 @@ func ParseArgumentsByExpression(key string, parameters
[]any) any {
return nil
}
var searchVal any
- _ = json.Unmarshal(data, &searchVal)
+ if err = json.Unmarshal(data, &searchVal); err != nil {
+ logger.Errorf("[Remoting][Polaris] unmarshal parameter %+v
failed, err=%v", parameters[index], err)
+ return nil
+ }
res, err := jsonpath.JsonPathLookup(searchVal, key)
if err != nil {
logger.Errorf("[Remoting][Polaris] invalid do json path lookup,
key=%s err=%v", key, err)