refactor: improve connection readiness check and enhance goroutine management in gRPC client; ensure proper context handling in stream listeners

This commit is contained in:
Marvin Zhang
2025-08-07 11:12:46 +08:00
parent 060396af3d
commit d042bc8cd7
3 changed files with 33 additions and 9 deletions

View File

@@ -162,13 +162,29 @@ func (sm *StreamManager) streamListener(ts *TaskStream) {
err error
}, 1)
// Start receive operation in a separate goroutine
// Start receive operation in a separate goroutine with proper cleanup
go func() {
defer func() {
if r := recover(); r != nil {
sm.service.Errorf("stream recv goroutine panic for task[%s]: %v", ts.taskId.Hex(), r)
}
}()
msg, err := ts.stream.Recv()
resultChan <- struct {
// Use select to ensure we don't block if the main goroutine has exited
select {
case resultChan <- struct {
msg *grpc.TaskServiceSubscribeResponse
err error
}{msg, err}
}{msg, err}:
case <-ts.ctx.Done():
// Parent context cancelled, just return without sending
return
case <-sm.ctx.Done():
// Manager context cancelled, just return without sending
return
}
}()
// Wait for result, timeout, or cancellation