-- JobStream (bidirectional) -----------------------------------------------
(stream listenerrpc.ListenerRPC_JobStreamServer)
| 198 | // -- JobStream (bidirectional) ----------------------------------------------- |
| 199 | |
| 200 | func (s *testServer) JobStream(stream listenerrpc.ListenerRPC_JobStreamServer) error { |
| 201 | s.mu.Lock() |
| 202 | missingListener := s.missingListener |
| 203 | disconnectCh := s.jobStreamDisconnectCh |
| 204 | s.mu.Unlock() |
| 205 | if missingListener { |
| 206 | return status.Error(codes.NotFound, "Listener not found") |
| 207 | } |
| 208 | |
| 209 | ctx := stream.Context() |
| 210 | errCh := make(chan error, 2) |
| 211 | |
| 212 | // Send job control messages to bridge. |
| 213 | go func() { |
| 214 | for { |
| 215 | select { |
| 216 | case ctrl, ok := <-s.jobCtrlCh: |
| 217 | if !ok { |
| 218 | return |
| 219 | } |
| 220 | if err := stream.Send(ctrl); err != nil { |
| 221 | errCh <- err |
| 222 | return |
| 223 | } |
| 224 | case <-ctx.Done(): |
| 225 | return |
| 226 | } |
| 227 | } |
| 228 | }() |
| 229 | |
| 230 | // Receive job status from bridge. |
| 231 | go func() { |
| 232 | for { |
| 233 | st, err := stream.Recv() |
| 234 | if err != nil { |
| 235 | errCh <- err |
| 236 | return |
| 237 | } |
| 238 | select { |
| 239 | case s.jobStatusCh <- st: |
| 240 | default: |
| 241 | } |
| 242 | } |
| 243 | }() |
| 244 | |
| 245 | select { |
| 246 | case <-ctx.Done(): |
| 247 | return ctx.Err() |
| 248 | case <-disconnectCh: |
| 249 | return status.Error(codes.Unavailable, "job stream disconnected") |
| 250 | case err := <-errCh: |
| 251 | return err |
| 252 | } |
| 253 | } |
| 254 | |
| 255 | // -- Helpers for reading captured state -------------------------------------- |
| 256 |
no test coverage detected