mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-05 22:12:04 +02:00
fix: drain regWait channel before closing to prevent message loss
- Add drain loop before closing regWait in reconnect() cleanup - Add drain loop before closing regWait in handleDisconnect() cleanup - Ensures no pending RegistrationResponse messages are lost during channel closure
This commit is contained in:
1 parent
2b184ed0b8
commit
6ad8cb56f7
1 file changed
+22
-4
+22
-4
@@ -270,8 +270,17 @@ func (c *GrpcAdminClient) reconnect(s *grpcState) error {
|
||||
c.safeCloseChannel(&s.streamExit)
|
||||
c.safeCloseChannel(&s.streamFailed)
|
||||
if s.regWait != nil {
|
||||
close(s.regWait)
|
||||
s.regWait = nil
|
||||
// Drain any pending registration responses before closing to avoid losing them
|
||||
for {
|
||||
select {
|
||||
case <-s.regWait:
|
||||
// continue draining until channel is empty
|
||||
default:
|
||||
close(s.regWait)
|
||||
s.regWait = nil
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
if s.streamCancel != nil {
|
||||
s.streamCancel()
|
||||
@@ -552,8 +561,17 @@ func (c *GrpcAdminClient) handleDisconnect(cmd grpcCommand, s *grpcState) {
|
||||
c.safeCloseChannel(&s.streamExit)
|
||||
c.safeCloseChannel(&s.streamFailed)
|
||||
if s.regWait != nil {
|
||||
close(s.regWait)
|
||||
s.regWait = nil
|
||||
// Drain any pending registration responses before closing to avoid losing them
|
||||
for {
|
||||
select {
|
||||
case <-s.regWait:
|
||||
// continue draining until channel is empty
|
||||
default:
|
||||
close(s.regWait)
|
||||
s.regWait = nil
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Cancel stream context
|
||||
|
||||
Reference in new issue
Block a user