| 145 | } |
| 146 | |
| 147 | func (s *speaker[S, R, _]) recvFromSerdes() { |
| 148 | defer close(s.recvLoopDone) |
| 149 | defer close(s.requests) |
| 150 | for { |
| 151 | select { |
| 152 | case <-s.ctx.Done(): |
| 153 | s.logger.Debug(s.ctx, "recvFromSerdes context done while waiting for proto", slog.Error(s.ctx.Err())) |
| 154 | return |
| 155 | case msg, ok := <-s.recvCh: |
| 156 | if !ok { |
| 157 | s.logger.Debug(s.ctx, "recvCh is closed") |
| 158 | return |
| 159 | } |
| 160 | rpc := msg.GetRpc() |
| 161 | if rpc != nil && rpc.ResponseTo != 0 { |
| 162 | // this is a unary response |
| 163 | s.tryToDeliverResponse(msg) |
| 164 | continue |
| 165 | } |
| 166 | req := &request[S, R]{ |
| 167 | ctx: s.ctx, |
| 168 | msg: msg, |
| 169 | replyCh: s.sendCh, |
| 170 | } |
| 171 | select { |
| 172 | case <-s.ctx.Done(): |
| 173 | s.logger.Debug(s.ctx, "recvFromSerdes context done while waiting for request handler", slog.Error(s.ctx.Err())) |
| 174 | return |
| 175 | case s.requests <- req: |
| 176 | } |
| 177 | } |
| 178 | } |
| 179 | } |
| 180 | |
| 181 | // Close closes the speaker |
| 182 | // nolint: revive |