rpc: use atomic types (#27214)

rpc: use atomic type
This commit is contained in:
s7v7nislands
2023-05-04 04:54:45 -04:00
committed by GitHub
parent dde2da0efb
commit ffda2c64c4
3 changed files with 9 additions and 9 deletions
+5 -5
View File
@@ -48,7 +48,7 @@ type Server struct {
mutex sync.Mutex
codecs map[ServerCodec]struct{}
run int32
run atomic.Bool
}
// NewServer creates a new server instance with no registered handlers.
@@ -56,8 +56,8 @@ func NewServer() *Server {
server := &Server{
idgen: randomIDGenerator(),
codecs: make(map[ServerCodec]struct{}),
run: 1,
}
server.run.Store(true)
// Register the default service providing meta information about the RPC service such
// as the services and methods it offers.
rpcService := &RPCService{server}
@@ -95,7 +95,7 @@ func (s *Server) trackCodec(codec ServerCodec) bool {
s.mutex.Lock()
defer s.mutex.Unlock()
if atomic.LoadInt32(&s.run) == 0 {
if !s.run.Load() {
return false // Don't serve if server is stopped.
}
s.codecs[codec] = struct{}{}
@@ -114,7 +114,7 @@ func (s *Server) untrackCodec(codec ServerCodec) {
// this mode.
func (s *Server) serveSingleRequest(ctx context.Context, codec ServerCodec) {
// Don't serve if server is stopped.
if atomic.LoadInt32(&s.run) == 0 {
if !s.run.Load() {
return
}
@@ -144,7 +144,7 @@ func (s *Server) Stop() {
s.mutex.Lock()
defer s.mutex.Unlock()
if atomic.CompareAndSwapInt32(&s.run, 1, 0) {
if s.run.CompareAndSwap(true, false) {
log.Debug("RPC server shutting down")
for codec := range s.codecs {
codec.close()