MCPcopy Create free account
hub / github.com/CyCoreSystems/ari-proxy / MultipleRequest

Method MultipleRequest

messagebus/nats.go:186–222  ·  view source on GitHub ↗

MultipleRequest sends a request message to multiple consumers

(topic string, req *proxy.Request, expectedResp int)

Source from the content-addressed store, hash-verified

184
185// MultipleRequest sends a request message to multiple consumers
186func (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
225func (n *NatsBus) MultipleRequestReturnFirstGoodResponse(topic string, req *proxy.Request, expectedResp int) (*proxy.Response, error) {

Callers

nothing calls this directly

Calls 4

NewMethod · 0.80
UnsubscribeMethod · 0.65
SubscribeMethod · 0.45
StopMethod · 0.45

Tested by

no test coverage detected