MultipleRequest sends a request message to multiple consumers
(topic string, req *proxy.Request, expectedResp int)
| 184 | |
| 185 | // MultipleRequest sends a request message to multiple consumers |
| 186 | func (n *NatsBus) MultipleRequest(topic string, req *proxy.Request, expectedResp int) ([]*proxy.Response, error) { |
| 187 | var responses []*proxy.Response |
| 188 | |
| 189 | reply := rid.New("rp") |
| 190 | |
| 191 | rf := &responseForwarder{ |
| 192 | expected: expectedResp, |
| 193 | fwdChan: make(chan *proxy.Response), |
| 194 | } |
| 195 | |
| 196 | replySub, err := n.conn.Subscribe(reply, rf.Forward) |
| 197 | if err != nil { |
| 198 | return nil, eris.Wrap(err, "failed to subscribe to data responses") |
| 199 | } |
| 200 | defer replySub.Unsubscribe() // nolint: errcheck |
| 201 | |
| 202 | // Make an all-call for the entity data |
| 203 | err = n.conn.PublishRequest(topic, reply, req) |
| 204 | if err != nil { |
| 205 | return nil, eris.Wrap(err, "failed to make request for data") |
| 206 | } |
| 207 | |
| 208 | // Wait for replies |
| 209 | timer := time.NewTimer(n.Config.RequestTimeout) |
| 210 | defer timer.Stop() |
| 211 | for { |
| 212 | select { |
| 213 | case <-timer.C: |
| 214 | return responses, nil |
| 215 | case resp, more := <-rf.fwdChan: |
| 216 | if !more { |
| 217 | return responses, nil |
| 218 | } |
| 219 | responses = append(responses, resp) |
| 220 | } |
| 221 | } |
| 222 | } |
| 223 | |
| 224 | // MultipleRequestReturnFirstGoodResponse sends a request message to multiple consumers and returns the first good response |
| 225 | func (n *NatsBus) MultipleRequestReturnFirstGoodResponse(topic string, req *proxy.Request, expectedResp int) (*proxy.Response, error) { |
nothing calls this directly
no test coverage detected