| 360 | } |
| 361 | } |
| 362 | TIBMsgHandle Send(TPtrArg<IIBPeer> peerArg, TRopeDataPacket* data, const TGUID& packetGuid) override { |
| 363 | TIBPeer* peer = static_cast<TIBPeer*>(peerArg.Get()); // trust me, I'm professional |
| 364 | if (peer == nullptr || peer->State != IIBPeer::OK) { |
| 365 | return -1; |
| 366 | } |
| 367 | Y_ASSERT(Channels.find(peer->QP->GetQPN())->second == peer); |
| 368 | TIBMsgHandle msgHandle = ++MsgCounter; |
| 369 | if (peer->SendCount >= MAX_SEND_COUNT) { |
| 370 | peer->PendingSendQueue.push_back(TPendingQueuedSend(packetGuid, msgHandle, data)); |
| 371 | } else { |
| 372 | //printf("Sending direct %d\n", msgHandle); |
| 373 | StartSend(peer, packetGuid, msgHandle, data); |
| 374 | } |
| 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); |
no test coverage detected