Request sends a proto message to the specific node
(ctx context.Context, id []byte, msg proto.Message)
| 318 | |
| 319 | // Request sends a proto message to the specific node |
| 320 | func (n *server) Request(ctx context.Context, id []byte, msg proto.Message) (p2pmsg P2PMessage, err error) { |
| 321 | var ( |
| 322 | result interface{} |
| 323 | ok bool |
| 324 | ) |
| 325 | addr := n.members.Lookup(id) |
| 326 | if addr == "" { |
| 327 | err = &P2PError{err: errors.New("No IP info"), dest: addr, t: time.Now()} |
| 328 | n.logger.Error(err) |
| 329 | return |
| 330 | } |
| 331 | defer n.logger.TimeTrack(time.Now(), "TimeRequest", nil) |
| 332 | opCtx, opCancel := context.WithTimeout(ctx, 5*time.Second) |
| 333 | defer opCancel() |
| 334 | req := NewP2pRequest(opCtx, sendReq, id, "", msg, 0) |
| 335 | if err = req.sendReq(n.calling); err != nil { |
| 336 | err = &P2PError{err: errors.Errorf("Request sendReq failed: %w", err), dest: addr, t: time.Now()} |
| 337 | n.logger.Error(err) |
| 338 | return |
| 339 | } |
| 340 | if result, err = req.waitForResult(); err != nil { |
| 341 | err = &P2PError{err: errors.Errorf("Request waitForResult: %w", err), dest: addr, t: time.Now()} |
| 342 | n.logger.Error(err) |
| 343 | return |
| 344 | } |
| 345 | |
| 346 | if p2pmsg, ok = result.(P2PMessage); !ok { |
| 347 | err = &P2PError{err: errors.Errorf("Request P2PMessage casting failed: %w", err), dest: addr, t: time.Now()} |
| 348 | n.logger.Error(err) |
| 349 | } |
| 350 | return |
| 351 | } |
| 352 | |
| 353 | // Reply sends a reply to the specific node |
| 354 | func (n *server) Reply(ctx context.Context, id []byte, nonce uint64, msg proto.Message) (err error) { |
nothing calls this directly
no test coverage detected