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