summaryrefslogtreecommitdiff
path: root/cpp/src/qpid/sys/rdma/RdmaIO.cpp
diff options
context:
space:
mode:
authorAndrew Stitcher <astitcher@apache.org>2010-12-23 17:11:09 +0000
committerAndrew Stitcher <astitcher@apache.org>2010-12-23 17:11:09 +0000
commit6a5298d1942290f17c8412f0312dd06067d4f835 (patch)
treea6709d1545f26ad68523b8e14b021df2fe25c5e8 /cpp/src/qpid/sys/rdma/RdmaIO.cpp
parentccda2c49458ad026e5c7a47bf4fb336a7e1fd6d5 (diff)
downloadqpid-python-6a5298d1942290f17c8412f0312dd06067d4f835.tar.gz
Factored rdma sending/receiving code out to make manipulating
credit isolated git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk/qpid@1052325 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'cpp/src/qpid/sys/rdma/RdmaIO.cpp')
-rw-r--r--cpp/src/qpid/sys/rdma/RdmaIO.cpp61
1 files changed, 38 insertions, 23 deletions
diff --git a/cpp/src/qpid/sys/rdma/RdmaIO.cpp b/cpp/src/qpid/sys/rdma/RdmaIO.cpp
index b47165a302..e082fdc416 100644
--- a/cpp/src/qpid/sys/rdma/RdmaIO.cpp
+++ b/cpp/src/qpid/sys/rdma/RdmaIO.cpp
@@ -172,18 +172,45 @@ namespace Rdma {
notifyCallback = nc;
}
+ void AsynchIO::queueBuffer(Buffer* buff, int credit) {
+ if (!buff) {
+ Buffer* ob = getBuffer();
+ // Have to send something as adapters hate it when you try to transfer 0 bytes
+ *reinterpret_cast< uint32_t* >(ob->bytes()) = htonl(credit);
+ ob->dataCount(sizeof(uint32_t));
+ qp->postSend(credit | IgnoreData, ob);
+ } else if (credit > 0) {
+ qp->postSend(credit, buff);
+ } else {
+ qp->postSend(buff);
+ }
+ }
+
+ Buffer* AsynchIO::extractBuffer(const QueuePairEvent& e) {
+ // Get our xmitCredit if it was sent
+ bool dataPresent = true;
+ if (e.immPresent() ) {
+ assert(xmitCredit>=0);
+ xmitCredit += (e.getImm() & ~FlagsMask);
+ dataPresent = ((e.getImm() & IgnoreData) == 0);
+ assert(xmitCredit>0);
+ }
+
+ Buffer* b = e.getBuffer();
+ if (!dataPresent) {
+ b->dataCount(0);
+ }
+ return b;
+ }
+
void AsynchIO::queueWrite(Buffer* buff) {
// Make sure we don't overrun our available buffers
// either at our end or the known available at the peers end
if (writable()) {
// TODO: We might want to batch up sending credit
- if (recvCredit > 0) {
- int creditSent = recvCredit & ~FlagsMask;
- qp->postSend(creditSent, buff);
- recvCredit -= creditSent;
- } else {
- qp->postSend(buff);
- }
+ int creditSent = recvCredit & ~FlagsMask;
+ queueBuffer(buff, creditSent);
+ recvCredit -= creditSent;
++outstandingWrites;
--xmitCredit;
assert(xmitCredit>=0);
@@ -315,22 +342,14 @@ namespace Rdma {
// Test if recv (or recv with imm)
//::ibv_wc_opcode eventType = e.getEventType();
- Buffer* b = e.getBuffer();
QueueDirection dir = e.getDirection();
if (dir == RECV) {
++recvEvents;
- // Get our xmitCredit if it was sent
- bool dataPresent = true;
- if (e.immPresent() ) {
- assert(xmitCredit>=0);
- xmitCredit += (e.getImm() & ~FlagsMask);
- dataPresent = ((e.getImm() & IgnoreData) == 0);
- assert(xmitCredit>0);
- }
+ Buffer* b = extractBuffer(e);
// if there was no data sent then the message was only to update our credit
- if ( dataPresent ) {
+ if ( b->dataCount() > 0 ) {
readCallback(*this, b);
}
@@ -347,13 +366,8 @@ namespace Rdma {
// but this is a little unlikely, as to get in this state we have to have received messages without sending any
// for a while so its likely we've received an credit update from the far side.
if (writable()) {
- Buffer* ob = getBuffer();
- // Have to send something as adapters hate it when you try to transfer 0 bytes
- *reinterpret_cast< uint32_t* >(ob->bytes()) = htonl(recvCredit);
- ob->dataCount(sizeof(uint32_t));
-
int creditSent = recvCredit & ~FlagsMask;
- qp->postSend(creditSent | IgnoreData, ob);
+ queueBuffer(0, creditSent);
recvCredit -= creditSent;
++outstandingWrites;
--xmitCredit;
@@ -363,6 +377,7 @@ namespace Rdma {
}
}
} else {
+ Buffer* b = e.getBuffer();
++sendEvents;
returnBuffer(b);
--outstandingWrites;