| 384 | return msgHandle; |
| 385 | } |
| 386 | void ParsePacket(ibv_wc* wc, NHPTimer::STime tCurrent) { |
| 387 | if (wc->status != IBV_WC_SUCCESS) { |
| 388 | TIBPeer* peer = GetChannelByQPN(wc->qp_num); |
| 389 | if (peer) { |
| 390 | //printf("failed recv packet (status %d)\n", wc->status); |
| 391 | PeerFailed(peer); |
| 392 | } else { |
| 393 | //printf("Ignoring recv error for closed/non existing QPN %d\n", wc->qp_num); |
| 394 | } |
| 395 | return; |
| 396 | } |
| 397 | |
| 398 | TIBRecvPacketProcess pkt(BP, *wc); |
| 399 | |
| 400 | TIBPeer* peer = GetChannelByQPN(wc->qp_num); |
| 401 | if (peer) { |
| 402 | Y_ASSERT(peer->State != IIBPeer::FAILED); |
| 403 | peer->LastRecv = tCurrent; |
| 404 | char cmdId = *(const char*)pkt.GetData(); |
| 405 | switch (cmdId) { |
| 406 | case CMD_CONFIRM: |
| 407 | //printf("got confirm\n"); |
| 408 | Y_ASSERT(peer->State == IIBPeer::CONNECTING); |
| 409 | peer->State = IIBPeer::OK; |
| 410 | break; |
| 411 | case CMD_DATA_TINY: |
| 412 | //printf("Recv CMD_DATA_TINY\n"); |
| 413 | { |
| 414 | const TCmdDataTiny& dataTiny = *(TCmdDataTiny*)pkt.GetData(); |
| 415 | TIBRequest* req = new TIBRequest; |
| 416 | req->ConnectionGuid = dataTiny.Header.ConnectionGuid; |
| 417 | req->Data = new TRopeDataPacket; |
| 418 | req->Data->Write(dataTiny.Data, dataTiny.Header.Size); |
| 419 | ReceivedList.push_back(req); |
| 420 | } |
| 421 | break; |
| 422 | case CMD_DATA_INIT: |
| 423 | //printf("Recv CMD_DATA_INIT\n"); |
| 424 | { |
| 425 | const TCmdDataInit& data = *(TCmdDataInit*)pkt.GetData(); |
| 426 | TIntrusivePtr<TIBMemBlock> blk = MemPool->Alloc(data.Size); |
| 427 | peer->RecvQueue.push_back(TQueuedRecv(data.PacketGuid, data.ConnectionGuid, blk)); |
| 428 | TCmdBufferReady ready; |
| 429 | ready.Command = CMD_BUFFER_READY; |
| 430 | ready.PacketGuid = data.PacketGuid; |
| 431 | |
| 432 | ready.RemoteAddr = reinterpret_cast<ui64>(blk->GetData()) / sizeof(char); |
| 433 | ready.RemoteKey = blk->GetMemRegion()->GetRKey(); |
| 434 | |
| 435 | peer->PostSend(BP, &ready, sizeof(ready), TCompleteInfo::CI_IGNORE, 0); |
| 436 | //printf("Send CMD_BUFFER_READY\n"); |
| 437 | } |
| 438 | break; |
| 439 | case CMD_BUFFER_READY: |
| 440 | //printf("Recv CMD_BUFFER_READY\n"); |
| 441 | { |
| 442 | const TCmdBufferReady& ready = *(TCmdBufferReady*)pkt.GetData(); |
| 443 | TDeque<TQueuedSend>::iterator z = peer->GetSend(ready.PacketGuid); |
nothing calls this directly
no test coverage detected