(stream ccintf.ChaincodeStream)
| 404 | } |
| 405 | |
| 406 | func (h *Handler) ProcessStream(stream ccintf.ChaincodeStream) error { |
| 407 | defer h.deregister() |
| 408 | |
| 409 | h.mutex.Lock() |
| 410 | h.streamDoneChan = make(chan struct{}) |
| 411 | h.mutex.Unlock() |
| 412 | defer close(h.streamDoneChan) |
| 413 | |
| 414 | h.chatStream = stream |
| 415 | h.errChan = make(chan error, 1) |
| 416 | |
| 417 | var keepaliveCh <-chan time.Time |
| 418 | if h.Keepalive != 0 { |
| 419 | ticker := time.NewTicker(h.Keepalive) |
| 420 | defer ticker.Stop() |
| 421 | keepaliveCh = ticker.C |
| 422 | } |
| 423 | |
| 424 | // holds return values from gRPC Recv below |
| 425 | type recvMsg struct { |
| 426 | msg *pb.ChaincodeMessage |
| 427 | err error |
| 428 | } |
| 429 | msgAvail := make(chan *recvMsg, 1) |
| 430 | |
| 431 | receiveMessage := func() { |
| 432 | in, err := h.chatStream.Recv() |
| 433 | msgAvail <- &recvMsg{in, err} |
| 434 | } |
| 435 | |
| 436 | go receiveMessage() |
| 437 | for { |
| 438 | select { |
| 439 | case rmsg := <-msgAvail: |
| 440 | switch { |
| 441 | // Defer the deregistering of the this handler. |
| 442 | case rmsg.err == io.EOF: |
| 443 | chaincodeLogger.Debugf("received EOF, ending chaincode support stream: %s", rmsg.err) |
| 444 | return rmsg.err |
| 445 | case rmsg.err != nil: |
| 446 | err := errors.Wrap(rmsg.err, "receive from chaincode support stream failed") |
| 447 | chaincodeLogger.Debugf("%+v", err) |
| 448 | return err |
| 449 | case rmsg.msg == nil: |
| 450 | err := errors.New("received nil message, ending chaincode support stream") |
| 451 | chaincodeLogger.Debugf("%+v", err) |
| 452 | return err |
| 453 | default: |
| 454 | err := h.handleMessage(rmsg.msg) |
| 455 | if err != nil { |
| 456 | err = errors.WithMessage(err, "error handling message, ending stream") |
| 457 | chaincodeLogger.Errorf("[%s] %+v", shorttxid(rmsg.msg.Txid), err) |
| 458 | return err |
| 459 | } |
| 460 | |
| 461 | go receiveMessage() |
| 462 | } |
| 463 |
no test coverage detected