commits
Threads by month
- ----- 2026 -----
- July
- June
- May
- April
- March
- February
- January
- ----- 2025 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2024 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2023 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2022 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2021 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2020 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2019 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2018 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2017 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2016 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2015 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2014 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2013 -----
- December
- November
- October
- September
- August
- July
- June
- May
- April
- March
- February
- January
- ----- 2012 -----
- December
- November
April 2015
- 1 participants
- 54 discussions
[mpich] MPICH primary repository branch, master, updated. v3.2b1-90-g51512b7
by noreply@mpich.org 21 Apr '15
by noreply@mpich.org 21 Apr '15
21 Apr '15
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "MPICH primary repository".
The branch, master has been updated
via 51512b7d5428d74c5d8ddae691b993596d570e7f (commit)
from fa10eeaa26b2188bbc63bf978a3ac8e4857e12df (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
http://git.mpich.org/mpich.git/commitdiff/51512b7d5428d74c5d8ddae691b993596…
commit 51512b7d5428d74c5d8ddae691b993596d570e7f
Author: Min Si <msi(a)il.is.s.u-tokyo.ac.jp>
Date: Tue Apr 21 18:38:15 2015 -0500
Revert "A workaround for FreeBSD pthread mallc/free bug."
This reverts commit 8f6b3cbb1f3a8a6f2202dbcff32fd76d7de7aa05.
diff --git a/test/mpi/threads/pt2pt/multisend4.c b/test/mpi/threads/pt2pt/multisend4.c
index 9b2f07b..041e6af 100644
--- a/test/mpi/threads/pt2pt/multisend4.c
+++ b/test/mpi/threads/pt2pt/multisend4.c
@@ -44,9 +44,9 @@ MTEST_THREAD_RETURN_TYPE run_test_sendrecv(void *arg)
fprintf( stderr, "Panic wsize = %d nthreads = %d\n",
wsize, nthreads );
- buf = (int *) malloc(2 * MAX_CNT * sizeof(int));
-
for (cnt=1; cnt < MAX_CNT; cnt = 2*cnt) {
+ buf = (int *)malloc( 2*cnt * sizeof(int) );
+
/* Wait for all senders to be ready */
MTest_thread_barrier(nthreads);
@@ -72,11 +72,10 @@ MTEST_THREAD_RETURN_TYPE run_test_sendrecv(void *arg)
t = MPI_Wtime() - t;
/* can't free the buffers until the requests are completed */
MTest_thread_barrier(nthreads);
+ free( buf );
if (thread_num == 1)
MTestPrintfMsg( 1, "buf size %d: time %f\n", cnt, t / MAX_LOOP );
}
-
- free(buf);
return (MTEST_THREAD_RETURN_TYPE)NULL;
}
-----------------------------------------------------------------------
Summary of changes:
test/mpi/threads/pt2pt/multisend4.c | 7 +++----
1 files changed, 3 insertions(+), 4 deletions(-)
hooks/post-receive
--
MPICH primary repository
1
0
[mpich] MPICH primary repository branch, master, updated. v3.2b1-89-gfa10eea
by noreply@mpich.org 21 Apr '15
by noreply@mpich.org 21 Apr '15
21 Apr '15
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "MPICH primary repository".
The branch, master has been updated
via fa10eeaa26b2188bbc63bf978a3ac8e4857e12df (commit)
from 8f6b3cbb1f3a8a6f2202dbcff32fd76d7de7aa05 (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
http://git.mpich.org/mpich.git/commitdiff/fa10eeaa26b2188bbc63bf978a3ac8e48…
commit fa10eeaa26b2188bbc63bf978a3ac8e4857e12df
Author: Charles J Archer <charles.j.archer(a)intel.com>
Date: Mon Apr 20 14:56:48 2015 -0700
Add ordering check to mprobe test case
Signed-off-by: Ken Raffenetti <raffenet(a)mcs.anl.gov>
diff --git a/test/mpi/pt2pt/mprobe.c b/test/mpi/pt2pt/mprobe.c
index f9abdc6..bc87d31 100644
--- a/test/mpi/pt2pt/mprobe.c
+++ b/test/mpi/pt2pt/mprobe.c
@@ -522,6 +522,62 @@ int main(int argc, char **argv)
}
MPI_Type_free(&vectype);
+ /* test 12: order test */
+ if (rank == 0) {
+ MPI_Request lrequest[2];
+ sendbuf[0] = 0xdeadbeef;
+ sendbuf[1] = 0xfeedface;
+ sendbuf[2] = 0xdeadbeef;
+ sendbuf[3] = 0xfeedface;
+ sendbuf[4] = 0xdeadbeef;
+ sendbuf[5] = 0xfeedface;
+ MPI_Isend(&sendbuf[0], 4, MPI_INT, 1, 6, MPI_COMM_WORLD,&lrequest[0]);
+ MPI_Isend(&sendbuf[4], 2, MPI_INT, 1, 6, MPI_COMM_WORLD,&lrequest[1]);
+ MPI_Waitall(2, &lrequest[0],MPI_STATUSES_IGNORE);
+ }
+ else {
+ memset(&s1, 0xab, sizeof(MPI_Status));
+ memset(&s2, 0xab, sizeof(MPI_Status));
+ /* the error field should remain unmodified */
+ s1.MPI_ERROR = MPI_ERR_DIMS;
+ s2.MPI_ERROR = MPI_ERR_TOPOLOGY;
+
+ msg = MPI_MESSAGE_NULL;
+ MPI_Mprobe(0, 6, MPI_COMM_WORLD, &msg, &s1);
+ check(s1.MPI_SOURCE == 0);
+ check(s1.MPI_TAG == 6);
+ check(s1.MPI_ERROR == MPI_ERR_DIMS);
+ check(msg != MPI_MESSAGE_NULL);
+
+ count = -1;
+ MPI_Get_count(&s1, MPI_INT, &count);
+ check(count == 4);
+
+ recvbuf[0] = 0x01234567;
+ recvbuf[1] = 0x89abcdef;
+ MPI_Recv(recvbuf, 2, MPI_INT, 0, 6, MPI_COMM_WORLD, &s2);
+ check(s2.MPI_SOURCE == 0);
+ check(s2.MPI_TAG == 6);
+ check(recvbuf[0] == 0xdeadbeef);
+ check(recvbuf[1] == 0xfeedface);
+
+ recvbuf[0] = 0x01234567;
+ recvbuf[1] = 0x89abcdef;
+ recvbuf[2] = 0x01234567;
+ recvbuf[3] = 0x89abcdef;
+ s2.MPI_ERROR = MPI_ERR_TOPOLOGY;
+
+ MPI_Mrecv(recvbuf, count, MPI_INT, &msg, &s2);
+ check(recvbuf[0] == 0xdeadbeef);
+ check(recvbuf[1] == 0xfeedface);
+ check(recvbuf[2] == 0xdeadbeef);
+ check(recvbuf[3] == 0xfeedface);
+ check(s2.MPI_SOURCE == 0);
+ check(s2.MPI_TAG == 6);
+ check(s2.MPI_ERROR == MPI_ERR_TOPOLOGY);
+ check(msg == MPI_MESSAGE_NULL);
+ }
+
free(sendbuf);
free(recvbuf);
-----------------------------------------------------------------------
Summary of changes:
test/mpi/pt2pt/mprobe.c | 56 +++++++++++++++++++++++++++++++++++++++++++++++
1 files changed, 56 insertions(+), 0 deletions(-)
hooks/post-receive
--
MPICH primary repository
1
0
[mpich] MPICH primary repository branch, master, updated. v3.2b1-88-g8f6b3cb
by noreply@mpich.org 20 Apr '15
by noreply@mpich.org 20 Apr '15
20 Apr '15
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "MPICH primary repository".
The branch, master has been updated
via 8f6b3cbb1f3a8a6f2202dbcff32fd76d7de7aa05 (commit)
from a56826750cc7b99150b9b59e602dcaf13cb5f661 (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
http://git.mpich.org/mpich.git/commitdiff/8f6b3cbb1f3a8a6f2202dbcff32fd76d7…
commit 8f6b3cbb1f3a8a6f2202dbcff32fd76d7de7aa05
Author: Min Si <msi(a)il.is.s.u-tokyo.ac.jp>
Date: Mon Apr 20 20:29:39 2015 -0500
A workaround for FreeBSD pthread mallc/free bug.
On FreeBSD, test threads/pt2pt/multisend4 sometimes reports the segfault
error when calling free function. This error only happens when the
buffer size is equal to 4M bytes and every thread performs malloc/free
for multiple times. This bug can be reproduced by using simple memcpy
without MPI communication, thus it is considered not as a MPI bug but a
bug of the thread-safe memory allocation on FreeBSD. A workaround of
this bug is to move malloc-free outside the loop to avoid frequent
malloc-free calls. This patch added it.
Signed-off-by: Huiwei Lu <huiweilu(a)mcs.anl.gov>
diff --git a/test/mpi/threads/pt2pt/multisend4.c b/test/mpi/threads/pt2pt/multisend4.c
index 041e6af..9b2f07b 100644
--- a/test/mpi/threads/pt2pt/multisend4.c
+++ b/test/mpi/threads/pt2pt/multisend4.c
@@ -44,9 +44,9 @@ MTEST_THREAD_RETURN_TYPE run_test_sendrecv(void *arg)
fprintf( stderr, "Panic wsize = %d nthreads = %d\n",
wsize, nthreads );
- for (cnt=1; cnt < MAX_CNT; cnt = 2*cnt) {
- buf = (int *)malloc( 2*cnt * sizeof(int) );
+ buf = (int *) malloc(2 * MAX_CNT * sizeof(int));
+ for (cnt=1; cnt < MAX_CNT; cnt = 2*cnt) {
/* Wait for all senders to be ready */
MTest_thread_barrier(nthreads);
@@ -72,10 +72,11 @@ MTEST_THREAD_RETURN_TYPE run_test_sendrecv(void *arg)
t = MPI_Wtime() - t;
/* can't free the buffers until the requests are completed */
MTest_thread_barrier(nthreads);
- free( buf );
if (thread_num == 1)
MTestPrintfMsg( 1, "buf size %d: time %f\n", cnt, t / MAX_LOOP );
}
+
+ free(buf);
return (MTEST_THREAD_RETURN_TYPE)NULL;
}
-----------------------------------------------------------------------
Summary of changes:
test/mpi/threads/pt2pt/multisend4.c | 7 ++++---
1 files changed, 4 insertions(+), 3 deletions(-)
hooks/post-receive
--
MPICH primary repository
1
0
[mpich] MPICH primary repository branch, master, updated. v3.2b1-87-ga568267
by noreply@mpich.org 20 Apr '15
by noreply@mpich.org 20 Apr '15
20 Apr '15
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "MPICH primary repository".
The branch, master has been updated
via a56826750cc7b99150b9b59e602dcaf13cb5f661 (commit)
via aec01b399f6ad205daf3fb5d6b3e52eb13cbe9c7 (commit)
via 001cd724bb0a6b06c5e26a47748f6b204ac0479f (commit)
via de0412c2980e5c8e34223ddceeafe8c006b9d66e (commit)
via 19f29078c001ff4176330815b066eac3da7b9e52 (commit)
from c09f396958cbef74684c3b423315b972973042a5 (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
http://git.mpich.org/mpich.git/commitdiff/a56826750cc7b99150b9b59e602dcaf13…
commit a56826750cc7b99150b9b59e602dcaf13cb5f661
Author: Xin Zhao <xinzhao3(a)illinois.edu>
Date: Mon Apr 20 15:41:07 2015 -0500
Modify comments about atomicity in GACC/FOP handlers.
Signed-off-by: Min Si <msi(a)il.is.s.u-tokyo.ac.jp>
Signed-off-by: Antonio J. Pena <apenya(a)mcs.anl.gov>
diff --git a/src/mpid/ch3/src/ch3u_handle_recv_req.c b/src/mpid/ch3/src/ch3u_handle_recv_req.c
index ff5ee54..5e62e66 100644
--- a/src/mpid/ch3/src/ch3u_handle_recv_req.c
+++ b/src/mpid/ch3/src/ch3u_handle_recv_req.c
@@ -263,7 +263,6 @@ int MPIDI_CH3_ReqHandler_GaccumRecvComplete(MPIDI_VC_t * vc, MPID_Request * rreq
MPID_Datatype_is_contig(rreq->dev.datatype, &is_contig);
- /* Copy data into a temporary buffer */
resp_req = MPID_Request_create();
MPIU_ERR_CHKANDJUMP(resp_req == NULL, mpi_errno, MPI_ERR_OTHER, "**nomemreq");
MPIU_Object_set_ref(resp_req, 1);
@@ -272,9 +271,13 @@ int MPIDI_CH3_ReqHandler_GaccumRecvComplete(MPIDI_VC_t * vc, MPID_Request * rreq
MPIU_CHKPMEM_MALLOC(resp_req->dev.user_buf, void *, rreq->dev.recv_data_sz,
mpi_errno, "GACC resp. buffer");
+ /* NOTE: 'copy data + ACC' needs to be atomic */
+
if (win_ptr->shm_allocated == TRUE)
MPIDI_CH3I_SHM_MUTEX_LOCK(win_ptr);
+ /* Copy data from target window to temporary buffer */
+
if (is_contig) {
MPIU_Memcpy(resp_req->dev.user_buf,
(void *) ((char *) rreq->dev.real_user_buf + rreq->dev.stream_offset),
@@ -397,6 +400,8 @@ int MPIDI_CH3_ReqHandler_FOPRecvComplete(MPIDI_VC_t * vc, MPID_Request * rreq, i
* operation are completed when counter reaches zero. */
win_ptr->at_completion_counter++;
+ /* NOTE: 'copy data + ACC' needs to be atomic */
+
if (win_ptr->shm_allocated == TRUE)
MPIDI_CH3I_SHM_MUTEX_LOCK(win_ptr);
@@ -1193,10 +1198,14 @@ static inline int perform_get_acc_in_lock_queue(MPID_Win * win_ptr,
get_accum_resp_pkt->flags |= MPIDI_CH3_PKT_FLAG_RMA_FLUSH_ACK;
get_accum_resp_pkt->target_rank = win_ptr->comm_ptr->rank;
- /* Perform ACCUMULATE OP */
+
+ /* NOTE: copy 'data + ACC' needs to be atomic */
+
if (win_ptr->shm_allocated == TRUE)
MPIDI_CH3I_SHM_MUTEX_LOCK(win_ptr);
+ /* Copy data from target window to response packet header */
+
void *src = (void *) (get_accum_pkt->addr), *dest =
(void *) (get_accum_resp_pkt->info.data);
mpi_errno = immed_copy(src, dest, len);
@@ -1206,6 +1215,8 @@ static inline int perform_get_acc_in_lock_queue(MPID_Win * win_ptr,
MPIU_ERR_POP(mpi_errno);
}
+ /* Perform ACCUMULATE OP */
+
/* All data fits in packet header */
/* NOTE: here we pass 0 as stream_offset to do_accumulate_op(), because the unit
that is piggybacked with LOCK flag must be the first stream unit */
@@ -1244,10 +1255,13 @@ static inline int perform_get_acc_in_lock_queue(MPID_Win * win_ptr,
MPID_Datatype_is_contig(get_accum_pkt->datatype, &is_contig);
- /* Perform ACCUMULATE OP */
+ /* NOTE: 'copy data + ACC' needs to be atomic */
+
if (win_ptr->shm_allocated == TRUE)
MPIDI_CH3I_SHM_MUTEX_LOCK(win_ptr);
+ /* Copy data from target window to temporary buffer */
+
/* NOTE: here we copy data from stream_offset = 0, because
the unit that is piggybacked with LOCK flag must be the
first stream unit. */
@@ -1271,6 +1285,8 @@ static inline int perform_get_acc_in_lock_queue(MPID_Win * win_ptr,
MPID_Segment_free(seg);
}
+ /* Perform ACCUMULATE OP */
+
/* NOTE: here we pass 0 as stream_offset to do_accumulate_op(), because the unit
that is piggybacked with LOCK flag must be the first stream unit */
mpi_errno = do_accumulate_op(lock_entry->data, get_accum_pkt->count, get_accum_pkt->datatype,
@@ -1379,9 +1395,13 @@ static inline int perform_fop_in_lock_queue(MPID_Win * win_ptr, MPIDI_RMA_Lock_e
win_ptr->at_completion_counter++;
}
+ /* NOTE: 'copy data + ACC' needs to be atomic */
+
if (win_ptr->shm_allocated == TRUE)
MPIDI_CH3I_SHM_MUTEX_LOCK(win_ptr);
+ /* Copy data from target window to temporary buffer / response packet header */
+
if (fop_pkt->type == MPIDI_CH3_PKT_FOP_IMMED) {
/* copy data to resp pkt header */
void *src = fop_pkt->addr, *dest = fop_resp_pkt->info.data;
diff --git a/src/mpid/ch3/src/ch3u_rma_pkthandler.c b/src/mpid/ch3/src/ch3u_rma_pkthandler.c
index 2bffebb..5ec61a4 100644
--- a/src/mpid/ch3/src/ch3u_rma_pkthandler.c
+++ b/src/mpid/ch3/src/ch3u_rma_pkthandler.c
@@ -836,6 +836,8 @@ int MPIDI_CH3_PktHandler_GetAccumulate(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
(get_accum_pkt->flags & MPIDI_CH3_PKT_FLAG_RMA_UNLOCK))
get_accum_resp_pkt->flags |= MPIDI_CH3_PKT_FLAG_RMA_FLUSH_ACK;
+ /* NOTE: 'copy data + ACC' needs to be atomic */
+
if (win_ptr->shm_allocated == TRUE)
MPIDI_CH3I_SHM_MUTEX_LOCK(win_ptr);
@@ -1232,6 +1234,8 @@ int MPIDI_CH3_PktHandler_FOP(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
(fop_pkt->flags & MPIDI_CH3_PKT_FLAG_RMA_UNLOCK))
fop_resp_pkt->flags |= MPIDI_CH3_PKT_FLAG_RMA_FLUSH_ACK;
+ /* NOTE: 'copy data + ACC' needs to be atomic */
+
if (win_ptr->shm_allocated == TRUE)
MPIDI_CH3I_SHM_MUTEX_LOCK(win_ptr);
http://git.mpich.org/mpich.git/commitdiff/aec01b399f6ad205daf3fb5d6b3e52eb1…
commit aec01b399f6ad205daf3fb5d6b3e52eb13cbe9c7
Author: Xin Zhao <xinzhao3(a)illinois.edu>
Date: Fri Apr 17 15:05:25 2015 -0500
Bug-fix: correct the wrong judgement in RMA function.
Here we should check if the packet type is FOP_IMMED, if so,
we initialize the response packet to be FOP_RESP_IMMED. The
originally code wrongly check the packet flag instead of packet
type.
Signed-off-by: Min Si <msi(a)il.is.s.u-tokyo.ac.jp>
Signed-off-by: Antonio J. Pena <apenya(a)mcs.anl.gov>
diff --git a/src/mpid/ch3/src/ch3u_handle_recv_req.c b/src/mpid/ch3/src/ch3u_handle_recv_req.c
index 7c85ba5..ff5ee54 100644
--- a/src/mpid/ch3/src/ch3u_handle_recv_req.c
+++ b/src/mpid/ch3/src/ch3u_handle_recv_req.c
@@ -1342,7 +1342,7 @@ static inline int perform_fop_in_lock_queue(MPID_Win * win_ptr, MPIDI_RMA_Lock_e
MPID_Datatype_is_contig(fop_pkt->datatype, &is_contig);
- if (fop_pkt->flags & MPIDI_CH3_PKT_FOP_IMMED) {
+ if (fop_pkt->type == MPIDI_CH3_PKT_FOP_IMMED) {
MPIDI_Pkt_init(fop_resp_pkt, MPIDI_CH3_PKT_FOP_RESP_IMMED);
}
else {
http://git.mpich.org/mpich.git/commitdiff/001cd724bb0a6b06c5e26a47748f6b204…
commit 001cd724bb0a6b06c5e26a47748f6b204ac0479f
Author: Xin Zhao <xinzhao3(a)illinois.edu>
Date: Fri Apr 17 13:14:40 2015 -0500
Delete assert that no longer makes sense.
After reducing the IMMED data size from 16 bytes to 8 bytes,
FOP data is no longer always fit in the packet header, hence
the assert no longer makes sense.
Signed-off-by: Min Si <msi(a)il.is.s.u-tokyo.ac.jp>
Signed-off-by: Antonio J. Pena <apenya(a)mcs.anl.gov>
diff --git a/src/mpid/ch3/src/ch3u_rma_pkthandler.c b/src/mpid/ch3/src/ch3u_rma_pkthandler.c
index 50800ee..2bffebb 100644
--- a/src/mpid/ch3/src/ch3u_rma_pkthandler.c
+++ b/src/mpid/ch3/src/ch3u_rma_pkthandler.c
@@ -1209,8 +1209,6 @@ int MPIDI_CH3_PktHandler_FOP(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
if (mpi_errno != MPI_SUCCESS)
MPIU_ERR_POP(mpi_errno);
- MPIU_Assert(rreq == NULL); /* FOP should not have request because all data
- * can fit in packet header */
if (acquire_lock_fail) {
(*rreqp) = rreq;
goto fn_exit;
http://git.mpich.org/mpich.git/commitdiff/de0412c2980e5c8e34223ddceeafe8c00…
commit de0412c2980e5c8e34223ddceeafe8c006b9d66e
Author: Xin Zhao <xinzhao3(a)illinois.edu>
Date: Fri Apr 17 09:04:44 2015 -0500
Set size of IMMED data in RMA packets to 8 bytes.
Originally the size of IMMED data in RMA packets is 16 bytes
which makes the size of CH3 packet be 56 bytes. Here we reduce
the size of IMMED data in RMA packets to 8 bytes, so that the
size of CH3 packet is reduced to 48 bytes, the same with
mpich-3.1.4 (the old RMA infrastructure).
Signed-off-by: Min Si <msi(a)il.is.s.u-tokyo.ac.jp>
Signed-off-by: Antonio J. Pena <apenya(a)mcs.anl.gov>
diff --git a/src/mpid/ch3/include/mpidpkt.h b/src/mpid/ch3/include/mpidpkt.h
index 0c4ff9e..eec2083 100644
--- a/src/mpid/ch3/include/mpidpkt.h
+++ b/src/mpid/ch3/include/mpidpkt.h
@@ -23,7 +23,7 @@
#define MPIDI_EAGER_SHORT_SIZE 16
/* This is the number of ints that can be carried within an RMA packet */
-#define MPIDI_RMA_IMMED_BYTES 16
+#define MPIDI_RMA_IMMED_BYTES 8
/* Union over all types (integer, logical, and multi-language types) that are
allowed in a CAS operation. This is used to allocate enough space in the
http://git.mpich.org/mpich.git/commitdiff/19f29078c001ff4176330815b066eac3d…
commit 19f29078c001ff4176330815b066eac3da7b9e52
Author: Xin Zhao <xinzhao3(a)illinois.edu>
Date: Sun Mar 22 14:54:28 2015 -0500
Move 'stream_offset' out of RMA packet struct.
'stream_offset' is used to specify the starting position
(on target window) of the current streaming unit in ACC-like
operations. It is originally put in the RMA packet struct,
which potentially increases the size of CH3 packet size.
In this patch, we move 'stream_offset' out of the RMA
packet as follows: 1. when target data is basic datatype,
we use 'stream_offset' and the starting address for the entire
operation to calculate the starting address for current
streaming unit, and rewrite 'addr' in RMA packet with that
value; 2. when target data is derived datatype, we cannot do
the same thing as basic datatype because the target needs to
know both the starting address for the entire operation and
the starting address for the current streaming unit. Therefore,
we send 'stream_offset' separately to the target side.
Signed-off-by: Min Si <msi(a)il.is.s.u-tokyo.ac.jp>
Signed-off-by: Antonio J. Pena <apenya(a)mcs.anl.gov>
diff --git a/src/mpid/ch3/include/mpid_rma_issue.h b/src/mpid/ch3/include/mpid_rma_issue.h
index bd40aac..17c6930 100644
--- a/src/mpid/ch3/include/mpid_rma_issue.h
+++ b/src/mpid/ch3/include/mpid_rma_issue.h
@@ -497,7 +497,7 @@ static int issue_from_origin_buffer_stream(MPIDI_RMA_Op_t * rma_op, MPIDI_VC_t *
MPID_Segment_pack_vector(segp, first, &last, dloop_vec, &vec_len);
- count = 2 + vec_len;
+ count = 3 + vec_len;
ints = (int *) MPIU_Malloc(sizeof(int) * (count + 1));
blocklens = &ints[1];
@@ -514,10 +514,16 @@ static int issue_from_origin_buffer_stream(MPIDI_RMA_Op_t * rma_op, MPIDI_VC_t *
MPIU_Assign_trunc(blocklens[1], target_dtp->dataloop_size, int);
datatypes[1] = MPI_BYTE;
+ req->dev.stream_offset = stream_offset;
+
+ displaces[2] = MPIU_PtrToAint(&(req->dev.stream_offset));
+ blocklens[2] = sizeof(req->dev.stream_offset);
+ datatypes[2] = MPI_BYTE;
+
for (i = 0; i < vec_len; i++) {
- displaces[i + 2] = MPIU_PtrToAint(dloop_vec[i].DLOOP_VECTOR_BUF);
- MPIU_Assign_trunc(blocklens[i + 2], dloop_vec[i].DLOOP_VECTOR_LEN, int);
- datatypes[i + 2] = MPI_BYTE;
+ displaces[i + 3] = MPIU_PtrToAint(dloop_vec[i].DLOOP_VECTOR_BUF);
+ MPIU_Assign_trunc(blocklens[i + 3], dloop_vec[i].DLOOP_VECTOR_LEN, int);
+ datatypes[i + 3] = MPI_BYTE;
}
MPID_Segment_free(segp);
@@ -730,7 +736,11 @@ static int issue_acc_op(MPIDI_RMA_Op_t * rma_op, MPID_Win * win_ptr,
stream_size = MPIR_MIN(stream_elem_count * predefined_dtp_size, rest_len);
rest_len -= stream_size;
- accum_pkt->info.metadata.stream_offset = stream_offset;
+ if (MPIR_DATATYPE_IS_PREDEFINED(accum_pkt->datatype)) {
+ accum_pkt->addr = (void *) ((char *) rma_op->original_target_addr
+ + j * stream_elem_count * predefined_dtp_extent);
+ accum_pkt->count = stream_size / predefined_dtp_size;
+ }
mpi_errno =
issue_from_origin_buffer_stream(rma_op, vc, stream_offset, stream_size, &curr_req);
@@ -939,7 +949,11 @@ static int issue_get_acc_op(MPIDI_RMA_Op_t * rma_op, MPID_Win * win_ptr,
stream_size = MPIR_MIN(stream_elem_count * predefined_dtp_size, rest_len);
rest_len -= stream_size;
- get_accum_pkt->info.metadata.stream_offset = stream_offset;
+ if (MPIR_DATATYPE_IS_PREDEFINED(get_accum_pkt->datatype)) {
+ get_accum_pkt->addr = (void *) ((char *) rma_op->original_target_addr
+ + j * stream_elem_count * predefined_dtp_extent);
+ get_accum_pkt->count = stream_size / predefined_dtp_size;
+ }
resp_req->dev.stream_offset = stream_offset;
diff --git a/src/mpid/ch3/include/mpid_rma_types.h b/src/mpid/ch3/include/mpid_rma_types.h
index be0e869..ff8c63d 100644
--- a/src/mpid/ch3/include/mpid_rma_types.h
+++ b/src/mpid/ch3/include/mpid_rma_types.h
@@ -60,6 +60,9 @@ typedef struct MPIDI_RMA_Op {
int result_count;
MPI_Datatype result_datatype;
+ /* used in streaming ACCs */
+ void *original_target_addr;
+
struct MPID_Request **reqs;
int reqs_size;
diff --git a/src/mpid/ch3/include/mpidpkt.h b/src/mpid/ch3/include/mpidpkt.h
index 3351a65..0c4ff9e 100644
--- a/src/mpid/ch3/include/mpidpkt.h
+++ b/src/mpid/ch3/include/mpidpkt.h
@@ -522,31 +522,16 @@ MPIDI_CH3_PKT_DEFS
err_ = MPI_SUCCESS; \
switch((pkt_).type) { \
case (MPIDI_CH3_PKT_PUT): \
- (pkt_).put.info.metadata.dataloop_size = (dataloop_size_); \
+ (pkt_).put.info.dataloop_size = (dataloop_size_); \
break; \
case (MPIDI_CH3_PKT_GET): \
- (pkt_).get.info.metadata.dataloop_size = (dataloop_size_); \
+ (pkt_).get.info.dataloop_size = (dataloop_size_); \
break; \
case (MPIDI_CH3_PKT_ACCUMULATE): \
- (pkt_).accum.info.metadata.dataloop_size = (dataloop_size_); \
+ (pkt_).accum.info.dataloop_size = (dataloop_size_); \
break; \
case (MPIDI_CH3_PKT_GET_ACCUM): \
- (pkt_).get_accum.info.metadata.dataloop_size = (dataloop_size_); \
- break; \
- default: \
- MPIU_ERR_SETANDJUMP1(err_, MPI_ERR_OTHER, "**invalidpkt", "**invalidpkt %d", (pkt_).type); \
- } \
- }
-
-#define MPIDI_CH3_PKT_RMA_GET_STREAM_OFFSET(pkt_, stream_offset_, err_) \
- { \
- err_ = MPI_SUCCESS; \
- switch((pkt_).type) { \
- case (MPIDI_CH3_PKT_ACCUMULATE): \
- (stream_offset_) = (pkt_).accum.info.metadata.stream_offset; \
- break; \
- case (MPIDI_CH3_PKT_GET_ACCUM): \
- (stream_offset_) = (pkt_).get_accum.info.metadata.stream_offset; \
+ (pkt_).get_accum.info.dataloop_size = (dataloop_size_); \
break; \
default: \
MPIU_ERR_SETANDJUMP1(err_, MPI_ERR_OTHER, "**invalidpkt", "**invalidpkt %d", (pkt_).type); \
@@ -609,12 +594,7 @@ typedef struct MPIDI_CH3_Pkt_put {
MPI_Win target_win_handle;
MPI_Win source_win_handle;
union {
- /* note that we use struct here in order
- * to consistently access dataloop_size
- * by "pkt->info.metadata.dataloop_size". */
- struct {
- int dataloop_size;
- } metadata;
+ int dataloop_size;
char data[MPIDI_RMA_IMMED_BYTES];
} info;
} MPIDI_CH3_Pkt_put_t;
@@ -628,10 +608,8 @@ typedef struct MPIDI_CH3_Pkt_get {
struct {
/* note that we use struct here in order
* to consistently access dataloop_size
- * by "pkt->info.metadata.dataloop_size". */
- struct {
- int dataloop_size; /* for derived datatypes */
- } metadata;
+ * by "pkt->info.dataloop_size". */
+ int dataloop_size; /* for derived datatypes */
} info;
MPI_Request request_handle;
MPI_Win target_win_handle;
@@ -662,10 +640,7 @@ typedef struct MPIDI_CH3_Pkt_accum {
MPI_Win target_win_handle;
MPI_Win source_win_handle;
union {
- struct {
- int dataloop_size;
- MPI_Aint stream_offset;
- } metadata;
+ int dataloop_size;
char data[MPIDI_RMA_IMMED_BYTES];
} info;
} MPIDI_CH3_Pkt_accum_t;
@@ -680,10 +655,7 @@ typedef struct MPIDI_CH3_Pkt_get_accum {
MPI_Op op;
MPI_Win target_win_handle;
union {
- struct {
- int dataloop_size;
- MPI_Aint stream_offset;
- } metadata;
+ int dataloop_size;
char data[MPIDI_RMA_IMMED_BYTES];
} info;
} MPIDI_CH3_Pkt_get_accum_t;
diff --git a/src/mpid/ch3/include/mpidrma.h b/src/mpid/ch3/include/mpidrma.h
index cc05214..4c90166 100644
--- a/src/mpid/ch3/include/mpidrma.h
+++ b/src/mpid/ch3/include/mpidrma.h
@@ -373,21 +373,8 @@ static inline int enqueue_lock_origin(MPID_Win * win_ptr, MPIDI_VC_t * vc,
MPID_Datatype_get_extent_macro(target_dtp, type_extent);
MPID_Datatype_get_size_macro(target_dtp, type_size);
- if (pkt->type == MPIDI_CH3_PKT_PUT) {
- recv_data_sz = type_size * target_count;
- buf_size = type_extent * target_count;
- }
- else {
- MPI_Aint stream_offset, stream_elem_count;
- MPI_Aint total_len, rest_len;
-
- MPIDI_CH3_PKT_RMA_GET_STREAM_OFFSET((*pkt), stream_offset, mpi_errno);
- stream_elem_count = MPIDI_CH3U_SRBuf_size / type_extent;
- total_len = type_size * target_count;
- rest_len = total_len - stream_offset;
- recv_data_sz = MPIR_MIN(rest_len, type_size * stream_elem_count);
- buf_size = type_extent * (recv_data_sz / type_size);
- }
+ recv_data_sz = type_size * target_count;
+ buf_size = type_extent * target_count;
if (new_ptr != NULL) {
if (win_ptr->current_lock_data_bytes + buf_size < MPIR_CVAR_CH3_RMA_LOCK_DATA_BYTES) {
diff --git a/src/mpid/ch3/src/ch3u_handle_recv_req.c b/src/mpid/ch3/src/ch3u_handle_recv_req.c
index 6e0d943..7c85ba5 100644
--- a/src/mpid/ch3/src/ch3u_handle_recv_req.c
+++ b/src/mpid/ch3/src/ch3u_handle_recv_req.c
@@ -1107,26 +1107,20 @@ static inline int perform_acc_in_lock_queue(MPID_Win * win_ptr, MPIDI_RMA_Lock_e
if (acc_pkt->type == MPIDI_CH3_PKT_ACCUMULATE_IMMED) {
/* All data fits in packet header */
+ /* NOTE: here we pass 0 as stream_offset to do_accumulate_op(), because the unit
+ that is piggybacked with LOCK flag must be the first stream unit */
mpi_errno = do_accumulate_op(acc_pkt->info.data, acc_pkt->count, acc_pkt->datatype,
acc_pkt->addr, acc_pkt->count, acc_pkt->datatype,
- 0, acc_pkt->op);
+ 0/* stream offset */, acc_pkt->op);
}
else {
MPIU_Assert(acc_pkt->type == MPIDI_CH3_PKT_ACCUMULATE);
- MPI_Aint type_size, type_extent;
- MPI_Aint total_len, rest_len, recv_count;
- MPID_Datatype_get_size_macro(acc_pkt->datatype, type_size);
- MPID_Datatype_get_extent_macro(acc_pkt->datatype, type_extent);
-
- total_len = type_size * acc_pkt->count;
- rest_len = total_len - acc_pkt->info.metadata.stream_offset;
- recv_count = MPIR_MIN((rest_len / type_size), (MPIDI_CH3U_SRBuf_size / type_extent));
- MPIU_Assert(recv_count > 0);
-
- mpi_errno = do_accumulate_op(lock_entry->data, recv_count, acc_pkt->datatype,
+ /* NOTE: here we pass 0 as stream_offset to do_accumulate_op(), because the unit
+ that is piggybacked with LOCK flag must be the first stream unit */
+ mpi_errno = do_accumulate_op(lock_entry->data, acc_pkt->count, acc_pkt->datatype,
acc_pkt->addr, acc_pkt->count, acc_pkt->datatype,
- acc_pkt->info.metadata.stream_offset, acc_pkt->op);
+ 0/* stream offset */, acc_pkt->op);
}
if (win_ptr->shm_allocated == TRUE)
@@ -1160,8 +1154,6 @@ static inline int perform_get_acc_in_lock_queue(MPID_Win * win_ptr,
MPID_IOV iov[MPID_IOV_LIMIT];
int is_contig;
int mpi_errno = MPI_SUCCESS;
- MPI_Aint type_extent;
- MPI_Aint total_len, rest_len, recv_count;
/* Piggyback candidate should have basic datatype for target datatype. */
MPIU_Assert(MPIR_DATATYPE_IS_PREDEFINED(get_accum_pkt->datatype));
@@ -1215,10 +1207,12 @@ static inline int perform_get_acc_in_lock_queue(MPID_Win * win_ptr,
}
/* All data fits in packet header */
+ /* NOTE: here we pass 0 as stream_offset to do_accumulate_op(), because the unit
+ that is piggybacked with LOCK flag must be the first stream unit */
mpi_errno =
do_accumulate_op(get_accum_pkt->info.data, get_accum_pkt->count,
get_accum_pkt->datatype, get_accum_pkt->addr, get_accum_pkt->count,
- get_accum_pkt->datatype, 0, get_accum_pkt->op);
+ get_accum_pkt->datatype, 0/* stream offset */, get_accum_pkt->op);
if (win_ptr->shm_allocated == TRUE)
MPIDI_CH3I_SHM_MUTEX_UNLOCK(win_ptr);
@@ -1246,14 +1240,7 @@ static inline int perform_get_acc_in_lock_queue(MPID_Win * win_ptr,
MPIU_Assert(get_accum_pkt->type == MPIDI_CH3_PKT_GET_ACCUM);
- MPID_Datatype_get_extent_macro(get_accum_pkt->datatype, type_extent);
-
- total_len = type_size * get_accum_pkt->count;
- rest_len = total_len - get_accum_pkt->info.metadata.stream_offset;
- recv_count = MPIR_MIN((rest_len / type_size), (MPIDI_CH3U_SRBuf_size / type_extent));
- MPIU_Assert(recv_count > 0);
-
- sreq->dev.user_buf = (void *) MPIU_Malloc(recv_count * type_size);
+ sreq->dev.user_buf = (void *) MPIU_Malloc(get_accum_pkt->count * type_size);
MPID_Datatype_is_contig(get_accum_pkt->datatype, &is_contig);
@@ -1261,15 +1248,16 @@ static inline int perform_get_acc_in_lock_queue(MPID_Win * win_ptr,
if (win_ptr->shm_allocated == TRUE)
MPIDI_CH3I_SHM_MUTEX_LOCK(win_ptr);
+ /* NOTE: here we copy data from stream_offset = 0, because
+ the unit that is piggybacked with LOCK flag must be the
+ first stream unit. */
if (is_contig) {
- MPIU_Memcpy(sreq->dev.user_buf,
- (void *) ((char *) get_accum_pkt->addr +
- get_accum_pkt->info.metadata.stream_offset), recv_count * type_size);
+ MPIU_Memcpy(sreq->dev.user_buf, get_accum_pkt->addr, get_accum_pkt->count * type_size);
}
else {
MPID_Segment *seg = MPID_Segment_alloc();
- MPI_Aint first = get_accum_pkt->info.metadata.stream_offset;
- MPI_Aint last = first + type_size * recv_count;
+ MPI_Aint first = 0;
+ MPI_Aint last = first + type_size * get_accum_pkt->count;
if (seg == NULL) {
if (win_ptr->shm_allocated == TRUE)
@@ -1283,9 +1271,11 @@ static inline int perform_get_acc_in_lock_queue(MPID_Win * win_ptr,
MPID_Segment_free(seg);
}
- mpi_errno = do_accumulate_op(lock_entry->data, recv_count, get_accum_pkt->datatype,
+ /* NOTE: here we pass 0 as stream_offset to do_accumulate_op(), because the unit
+ that is piggybacked with LOCK flag must be the first stream unit */
+ mpi_errno = do_accumulate_op(lock_entry->data, get_accum_pkt->count, get_accum_pkt->datatype,
get_accum_pkt->addr, get_accum_pkt->count, get_accum_pkt->datatype,
- get_accum_pkt->info.metadata.stream_offset, get_accum_pkt->op);
+ 0/* stream offset */, get_accum_pkt->op);
if (win_ptr->shm_allocated == TRUE)
MPIDI_CH3I_SHM_MUTEX_UNLOCK(win_ptr);
@@ -1311,7 +1301,7 @@ static inline int perform_get_acc_in_lock_queue(MPID_Win * win_ptr,
iov[0].MPID_IOV_BUF = (MPID_IOV_BUF_CAST) get_accum_resp_pkt;
iov[0].MPID_IOV_LEN = sizeof(*get_accum_resp_pkt);
iov[1].MPID_IOV_BUF = (MPID_IOV_BUF_CAST) ((char *) sreq->dev.user_buf);
- iov[1].MPID_IOV_LEN = recv_count * type_size;
+ iov[1].MPID_IOV_LEN = get_accum_pkt->count * type_size;
iovcnt = 2;
mpi_errno = MPIDI_CH3_iSendv(lock_entry->vc, sreq, iov, iovcnt);
diff --git a/src/mpid/ch3/src/ch3u_rma_ops.c b/src/mpid/ch3/src/ch3u_rma_ops.c
index 9611716..136c260 100644
--- a/src/mpid/ch3/src/ch3u_rma_ops.c
+++ b/src/mpid/ch3/src/ch3u_rma_ops.c
@@ -198,7 +198,7 @@ int MPIDI_CH3I_Put(const void *origin_addr, int origin_count, MPI_Datatype
win_ptr->basic_info_table[target_rank].disp_unit * target_disp;
put_pkt->count = target_count;
put_pkt->datatype = target_datatype;
- put_pkt->info.metadata.dataloop_size = 0;
+ put_pkt->info.dataloop_size = 0;
put_pkt->target_win_handle = win_ptr->basic_info_table[target_rank].win_handle;
put_pkt->source_win_handle = win_ptr->handle;
put_pkt->flags = MPIDI_CH3_PKT_FLAG_NONE;
@@ -387,7 +387,7 @@ int MPIDI_CH3I_Get(void *origin_addr, int origin_count, MPI_Datatype
win_ptr->basic_info_table[target_rank].disp_unit * target_disp;
get_pkt->count = target_count;
get_pkt->datatype = target_datatype;
- get_pkt->info.metadata.dataloop_size = 0;
+ get_pkt->info.dataloop_size = 0;
get_pkt->target_win_handle = win_ptr->basic_info_table[target_rank].win_handle;
get_pkt->flags = MPIDI_CH3_PKT_FLAG_NONE;
if (use_immed_resp_pkt)
@@ -612,11 +612,10 @@ int MPIDI_CH3I_Accumulate(const void *origin_addr, int origin_count, MPI_Datatyp
win_ptr->basic_info_table[target_rank].disp_unit * target_disp;
accum_pkt->count = target_count;
accum_pkt->datatype = target_datatype;
- accum_pkt->info.metadata.dataloop_size = 0;
+ accum_pkt->info.dataloop_size = 0;
accum_pkt->op = op;
accum_pkt->target_win_handle = win_ptr->basic_info_table[target_rank].win_handle;
accum_pkt->source_win_handle = win_ptr->handle;
- accum_pkt->info.metadata.stream_offset = 0;
accum_pkt->flags = MPIDI_CH3_PKT_FLAG_NONE;
if (use_immed_pkt) {
void *src = (void *) origin_addr, *dest = (void *) (accum_pkt->info.data);
@@ -625,6 +624,12 @@ int MPIDI_CH3I_Accumulate(const void *origin_addr, int origin_count, MPI_Datatyp
MPIU_ERR_POP(mpi_errno);
}
+ /* NOTE: here we backup the original starting address for the entire operation
+ on target window in 'original_target_addr', because when actually issuing
+ this operation, we may stream this operation and overwrite 'addr' with the
+ starting address for the streaming unit. */
+ new_ptr->original_target_addr = accum_pkt->addr;
+
MPIR_T_PVAR_TIMER_END(RMA, rma_rmaqueue_set);
mpi_errno = MPIDI_CH3I_Win_enqueue_op(win_ptr, new_ptr);
@@ -810,7 +815,7 @@ int MPIDI_CH3I_Get_accumulate(const void *origin_addr, int origin_count,
win_ptr->basic_info_table[target_rank].disp_unit * target_disp;
get_pkt->count = target_count;
get_pkt->datatype = target_datatype;
- get_pkt->info.metadata.dataloop_size = 0;
+ get_pkt->info.dataloop_size = 0;
get_pkt->target_win_handle = win_ptr->basic_info_table[target_rank].win_handle;
get_pkt->flags = MPIDI_CH3_PKT_FLAG_NONE;
if (use_immed_resp_pkt == TRUE)
@@ -931,10 +936,9 @@ int MPIDI_CH3I_Get_accumulate(const void *origin_addr, int origin_count,
win_ptr->basic_info_table[target_rank].disp_unit * target_disp;
get_accum_pkt->count = target_count;
get_accum_pkt->datatype = target_datatype;
- get_accum_pkt->info.metadata.dataloop_size = 0;
+ get_accum_pkt->info.dataloop_size = 0;
get_accum_pkt->op = op;
get_accum_pkt->target_win_handle = win_ptr->basic_info_table[target_rank].win_handle;
- get_accum_pkt->info.metadata.stream_offset = 0;
get_accum_pkt->flags = MPIDI_CH3_PKT_FLAG_NONE;
if (use_immed_pkt) {
void *src = (void *) origin_addr, *dest = (void *) (get_accum_pkt->info.data);
@@ -942,6 +946,12 @@ int MPIDI_CH3I_Get_accumulate(const void *origin_addr, int origin_count,
if (mpi_errno != MPI_SUCCESS)
MPIU_ERR_POP(mpi_errno);
}
+
+ /* NOTE: here we backup the original starting address for the entire operation
+ on target window in 'original_target_addr', because when actually issuing
+ this operation, we may stream this operation and overwrite 'addr' with the
+ starting address for the streaming unit. */
+ new_ptr->original_target_addr = get_accum_pkt->addr;
}
MPIR_T_PVAR_TIMER_END(RMA, rma_rmaqueue_set);
@@ -1346,7 +1356,7 @@ int MPIDI_Fetch_and_op(const void *origin_addr, void *result_addr,
win_ptr->basic_info_table[target_rank].disp_unit * target_disp;
get_pkt->count = 1;
get_pkt->datatype = datatype;
- get_pkt->info.metadata.dataloop_size = 0;
+ get_pkt->info.dataloop_size = 0;
get_pkt->target_win_handle = win_ptr->basic_info_table[target_rank].win_handle;
get_pkt->flags = MPIDI_CH3_PKT_FLAG_NONE;
if (use_immed_resp_pkt == TRUE)
diff --git a/src/mpid/ch3/src/ch3u_rma_pkthandler.c b/src/mpid/ch3/src/ch3u_rma_pkthandler.c
index adfe90b..50800ee 100644
--- a/src/mpid/ch3/src/ch3u_rma_pkthandler.c
+++ b/src/mpid/ch3/src/ch3u_rma_pkthandler.c
@@ -293,24 +293,24 @@ int MPIDI_CH3_PktHandler_Put(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
"MPIDI_RMA_dtype_info");
}
- req->dev.dataloop = MPIU_Malloc(put_pkt->info.metadata.dataloop_size);
+ req->dev.dataloop = MPIU_Malloc(put_pkt->info.dataloop_size);
if (!req->dev.dataloop) {
MPIU_ERR_SETANDJUMP1(mpi_errno, MPI_ERR_OTHER, "**nomem", "**nomem %d",
- put_pkt->info.metadata.dataloop_size);
+ put_pkt->info.dataloop_size);
}
/* if we received all of the dtype_info and dataloop, copy it
* now and call the handler, otherwise set the iov and let the
* channel copy it */
- if (data_len >= sizeof(MPIDI_RMA_dtype_info) + put_pkt->info.metadata.dataloop_size) {
+ if (data_len >= sizeof(MPIDI_RMA_dtype_info) + put_pkt->info.dataloop_size) {
/* copy all of dtype_info and dataloop */
MPIU_Memcpy(req->dev.dtype_info, data_buf, sizeof(MPIDI_RMA_dtype_info));
MPIU_Memcpy(req->dev.dataloop, data_buf + sizeof(MPIDI_RMA_dtype_info),
- put_pkt->info.metadata.dataloop_size);
+ put_pkt->info.dataloop_size);
*buflen =
sizeof(MPIDI_CH3_Pkt_t) + sizeof(MPIDI_RMA_dtype_info) +
- put_pkt->info.metadata.dataloop_size;
+ put_pkt->info.dataloop_size;
/* All dtype data has been received, call req handler */
mpi_errno = MPIDI_CH3_ReqHandler_PutDerivedDTRecvComplete(vc, req, &complete);
@@ -325,7 +325,7 @@ int MPIDI_CH3_PktHandler_Put(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
req->dev.iov[0].MPID_IOV_BUF = (MPID_IOV_BUF_CAST) ((char *) req->dev.dtype_info);
req->dev.iov[0].MPID_IOV_LEN = sizeof(MPIDI_RMA_dtype_info);
req->dev.iov[1].MPID_IOV_BUF = (MPID_IOV_BUF_CAST) req->dev.dataloop;
- req->dev.iov[1].MPID_IOV_LEN = put_pkt->info.metadata.dataloop_size;
+ req->dev.iov[1].MPID_IOV_LEN = put_pkt->info.dataloop_size;
req->dev.iov_count = 2;
*buflen = sizeof(MPIDI_CH3_Pkt_t);
@@ -519,24 +519,24 @@ int MPIDI_CH3_PktHandler_Get(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
"MPIDI_RMA_dtype_info");
}
- req->dev.dataloop = MPIU_Malloc(get_pkt->info.metadata.dataloop_size);
+ req->dev.dataloop = MPIU_Malloc(get_pkt->info.dataloop_size);
if (!req->dev.dataloop) {
MPIU_ERR_SETANDJUMP1(mpi_errno, MPI_ERR_OTHER, "**nomem", "**nomem %d",
- get_pkt->info.metadata.dataloop_size);
+ get_pkt->info.dataloop_size);
}
/* if we received all of the dtype_info and dataloop, copy it
* now and call the handler, otherwise set the iov and let the
* channel copy it */
- if (data_len >= sizeof(MPIDI_RMA_dtype_info) + get_pkt->info.metadata.dataloop_size) {
+ if (data_len >= sizeof(MPIDI_RMA_dtype_info) + get_pkt->info.dataloop_size) {
/* copy all of dtype_info and dataloop */
MPIU_Memcpy(req->dev.dtype_info, data_buf, sizeof(MPIDI_RMA_dtype_info));
MPIU_Memcpy(req->dev.dataloop, data_buf + sizeof(MPIDI_RMA_dtype_info),
- get_pkt->info.metadata.dataloop_size);
+ get_pkt->info.dataloop_size);
*buflen =
sizeof(MPIDI_CH3_Pkt_t) + sizeof(MPIDI_RMA_dtype_info) +
- get_pkt->info.metadata.dataloop_size;
+ get_pkt->info.dataloop_size;
/* All dtype data has been received, call req handler */
mpi_errno = MPIDI_CH3_ReqHandler_GetDerivedDTRecvComplete(vc, req, &complete);
@@ -549,7 +549,7 @@ int MPIDI_CH3_PktHandler_Get(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
req->dev.iov[0].MPID_IOV_BUF = (MPID_IOV_BUF_CAST) req->dev.dtype_info;
req->dev.iov[0].MPID_IOV_LEN = sizeof(MPIDI_RMA_dtype_info);
req->dev.iov[1].MPID_IOV_BUF = (MPID_IOV_BUF_CAST) req->dev.dataloop;
- req->dev.iov[1].MPID_IOV_LEN = get_pkt->info.metadata.dataloop_size;
+ req->dev.iov[1].MPID_IOV_LEN = get_pkt->info.dataloop_size;
req->dev.iov_count = 2;
*buflen = sizeof(MPIDI_CH3_Pkt_t);
@@ -574,7 +574,6 @@ int MPIDI_CH3_PktHandler_Accumulate(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
{
MPIDI_CH3_Pkt_accum_t *accum_pkt = &pkt->accum;
MPID_Request *req = NULL;
- MPI_Aint extent;
int complete = 0;
char *data_buf = NULL;
MPIDI_msg_sz_t data_len;
@@ -582,7 +581,6 @@ int MPIDI_CH3_PktHandler_Accumulate(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
int acquire_lock_fail = 0;
int mpi_errno = MPI_SUCCESS;
MPI_Aint type_size;
- MPI_Aint stream_elem_count, rest_len, total_len;
MPIDI_STATE_DECL(MPID_STATE_MPIDI_CH3_PKTHANDLER_ACCUMULATE);
MPIDI_FUNC_ENTER(MPID_STATE_MPIDI_CH3_PKTHANDLER_ACCUMULATE);
@@ -640,7 +638,6 @@ int MPIDI_CH3_PktHandler_Accumulate(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
req->dev.target_win_handle = accum_pkt->target_win_handle;
req->dev.source_win_handle = accum_pkt->source_win_handle;
req->dev.flags = accum_pkt->flags;
- req->dev.stream_offset = accum_pkt->info.metadata.stream_offset;
req->dev.resp_request_handle = MPI_REQUEST_NULL;
req->dev.OnFinal = MPIDI_CH3_ReqHandler_AccumRecvComplete;
@@ -653,8 +650,6 @@ int MPIDI_CH3_PktHandler_Accumulate(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
MPIDI_Request_set_type(req, MPIDI_REQUEST_TYPE_ACCUM_RECV);
req->dev.datatype = accum_pkt->datatype;
- MPID_Datatype_get_extent_macro(accum_pkt->datatype, extent);
-
MPIU_Assert(!MPIDI_Request_get_srbuf_flag(req));
/* allocate a SRBuf for receiving stream unit */
MPIDI_CH3U_SRBuf_alloc(req, MPIDI_CH3U_SRBuf_size);
@@ -673,11 +668,7 @@ int MPIDI_CH3_PktHandler_Accumulate(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
MPID_Datatype_get_size_macro(accum_pkt->datatype, type_size);
- total_len = type_size * accum_pkt->count;
- rest_len = total_len - req->dev.stream_offset;
- stream_elem_count = MPIDI_CH3U_SRBuf_size / extent;
-
- req->dev.recv_data_sz = MPIR_MIN(rest_len, stream_elem_count * type_size);
+ req->dev.recv_data_sz = type_size * accum_pkt->count;
MPIU_Assert(req->dev.recv_data_sz > 0);
mpi_errno = MPIDI_CH3U_Receive_data_found(req, data_buf, &data_len, &complete);
@@ -709,21 +700,25 @@ int MPIDI_CH3_PktHandler_Accumulate(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
"MPIDI_RMA_dtype_info");
}
- req->dev.dataloop = MPIU_Malloc(accum_pkt->info.metadata.dataloop_size);
+ req->dev.dataloop = MPIU_Malloc(accum_pkt->info.dataloop_size);
if (!req->dev.dataloop) {
MPIU_ERR_SETANDJUMP1(mpi_errno, MPI_ERR_OTHER, "**nomem", "**nomem %d",
- accum_pkt->info.metadata.dataloop_size);
+ accum_pkt->info.dataloop_size);
}
- if (data_len >= sizeof(MPIDI_RMA_dtype_info) + accum_pkt->info.metadata.dataloop_size) {
+ if (data_len >= sizeof(MPIDI_RMA_dtype_info) + accum_pkt->info.dataloop_size +
+ sizeof(req->dev.stream_offset)) {
/* copy all of dtype_info and dataloop */
MPIU_Memcpy(req->dev.dtype_info, data_buf, sizeof(MPIDI_RMA_dtype_info));
MPIU_Memcpy(req->dev.dataloop, data_buf + sizeof(MPIDI_RMA_dtype_info),
- accum_pkt->info.metadata.dataloop_size);
+ accum_pkt->info.dataloop_size);
+ MPIU_Memcpy(&(req->dev.stream_offset),
+ data_buf + sizeof(MPIDI_RMA_dtype_info) + accum_pkt->info.dataloop_size,
+ sizeof(req->dev.stream_offset));
*buflen =
sizeof(MPIDI_CH3_Pkt_t) + sizeof(MPIDI_RMA_dtype_info) +
- accum_pkt->info.metadata.dataloop_size;
+ accum_pkt->info.dataloop_size + sizeof(req->dev.stream_offset);
/* All dtype data has been received, call req handler */
mpi_errno = MPIDI_CH3_ReqHandler_AccumDerivedDTRecvComplete(vc, req, &complete);
@@ -738,8 +733,10 @@ int MPIDI_CH3_PktHandler_Accumulate(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
req->dev.iov[0].MPID_IOV_BUF = (MPID_IOV_BUF_CAST) req->dev.dtype_info;
req->dev.iov[0].MPID_IOV_LEN = sizeof(MPIDI_RMA_dtype_info);
req->dev.iov[1].MPID_IOV_BUF = (MPID_IOV_BUF_CAST) req->dev.dataloop;
- req->dev.iov[1].MPID_IOV_LEN = accum_pkt->info.metadata.dataloop_size;
- req->dev.iov_count = 2;
+ req->dev.iov[1].MPID_IOV_LEN = accum_pkt->info.dataloop_size;
+ req->dev.iov[2].MPID_IOV_BUF = &(req->dev.stream_offset);
+ req->dev.iov[2].MPID_IOV_LEN = sizeof(req->dev.stream_offset);
+ req->dev.iov_count = 3;
*buflen = sizeof(MPIDI_CH3_Pkt_t);
}
@@ -769,14 +766,12 @@ int MPIDI_CH3_PktHandler_GetAccumulate(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
{
MPIDI_CH3_Pkt_get_accum_t *get_accum_pkt = &pkt->get_accum;
MPID_Request *req = NULL;
- MPI_Aint extent;
int complete = 0;
char *data_buf = NULL;
MPIDI_msg_sz_t data_len;
MPID_Win *win_ptr;
int acquire_lock_fail = 0;
int mpi_errno = MPI_SUCCESS;
- MPI_Aint stream_elem_count, rest_len, total_len;
MPIDI_STATE_DECL(MPID_STATE_MPIDI_CH3_PKTHANDLER_GETACCUMULATE);
MPIDI_FUNC_ENTER(MPID_STATE_MPIDI_CH3_PKTHANDLER_GETACCUMULATE);
@@ -892,7 +887,6 @@ int MPIDI_CH3_PktHandler_GetAccumulate(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
req->dev.real_user_buf = get_accum_pkt->addr;
req->dev.target_win_handle = get_accum_pkt->target_win_handle;
req->dev.flags = get_accum_pkt->flags;
- req->dev.stream_offset = get_accum_pkt->info.metadata.stream_offset;
req->dev.resp_request_handle = get_accum_pkt->request_handle;
req->dev.OnFinal = MPIDI_CH3_ReqHandler_GaccumRecvComplete;
@@ -907,8 +901,6 @@ int MPIDI_CH3_PktHandler_GetAccumulate(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
MPIDI_Request_set_type(req, MPIDI_REQUEST_TYPE_GET_ACCUM_RECV);
req->dev.datatype = get_accum_pkt->datatype;
- MPID_Datatype_get_extent_macro(get_accum_pkt->datatype, extent);
-
MPIU_Assert(!MPIDI_Request_get_srbuf_flag(req));
/* allocate a SRBuf for receiving stream unit */
MPIDI_CH3U_SRBuf_alloc(req, MPIDI_CH3U_SRBuf_size);
@@ -926,11 +918,7 @@ int MPIDI_CH3_PktHandler_GetAccumulate(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
req->dev.user_buf = req->dev.tmpbuf;
MPID_Datatype_get_size_macro(get_accum_pkt->datatype, type_size);
- total_len = type_size * get_accum_pkt->count;
- rest_len = total_len - req->dev.stream_offset;
- stream_elem_count = MPIDI_CH3U_SRBuf_size / extent;
-
- req->dev.recv_data_sz = MPIR_MIN(rest_len, stream_elem_count * type_size);
+ req->dev.recv_data_sz = type_size * get_accum_pkt->count;
MPIU_Assert(req->dev.recv_data_sz > 0);
mpi_errno = MPIDI_CH3U_Receive_data_found(req, data_buf, &data_len, &complete);
@@ -962,22 +950,26 @@ int MPIDI_CH3_PktHandler_GetAccumulate(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
"MPIDI_RMA_dtype_info");
}
- req->dev.dataloop = MPIU_Malloc(get_accum_pkt->info.metadata.dataloop_size);
+ req->dev.dataloop = MPIU_Malloc(get_accum_pkt->info.dataloop_size);
if (!req->dev.dataloop) {
MPIU_ERR_SETANDJUMP1(mpi_errno, MPI_ERR_OTHER, "**nomem", "**nomem %d",
- get_accum_pkt->info.metadata.dataloop_size);
+ get_accum_pkt->info.dataloop_size);
}
if (data_len >=
- sizeof(MPIDI_RMA_dtype_info) + get_accum_pkt->info.metadata.dataloop_size) {
+ sizeof(MPIDI_RMA_dtype_info) + get_accum_pkt->info.dataloop_size +
+ sizeof(req->dev.stream_offset)) {
/* copy all of dtype_info and dataloop */
MPIU_Memcpy(req->dev.dtype_info, data_buf, sizeof(MPIDI_RMA_dtype_info));
MPIU_Memcpy(req->dev.dataloop, data_buf + sizeof(MPIDI_RMA_dtype_info),
- get_accum_pkt->info.metadata.dataloop_size);
+ get_accum_pkt->info.dataloop_size);
+ MPIU_Memcpy(&(req->dev.stream_offset),
+ data_buf + sizeof(MPIDI_RMA_dtype_info) +
+ get_accum_pkt->info.dataloop_size, sizeof(req->dev.stream_offset));
*buflen =
sizeof(MPIDI_CH3_Pkt_t) + sizeof(MPIDI_RMA_dtype_info) +
- get_accum_pkt->info.metadata.dataloop_size;
+ get_accum_pkt->info.dataloop_size + sizeof(req->dev.stream_offset);
/* All dtype data has been received, call req handler */
mpi_errno = MPIDI_CH3_ReqHandler_GaccumDerivedDTRecvComplete(vc, req, &complete);
@@ -992,8 +984,10 @@ int MPIDI_CH3_PktHandler_GetAccumulate(MPIDI_VC_t * vc, MPIDI_CH3_Pkt_t * pkt,
req->dev.iov[0].MPID_IOV_BUF = (MPID_IOV_BUF_CAST) req->dev.dtype_info;
req->dev.iov[0].MPID_IOV_LEN = sizeof(MPIDI_RMA_dtype_info);
req->dev.iov[1].MPID_IOV_BUF = (MPID_IOV_BUF_CAST) req->dev.dataloop;
- req->dev.iov[1].MPID_IOV_LEN = get_accum_pkt->info.metadata.dataloop_size;
- req->dev.iov_count = 2;
+ req->dev.iov[1].MPID_IOV_LEN = get_accum_pkt->info.dataloop_size;
+ req->dev.iov[2].MPID_IOV_BUF = &(req->dev.stream_offset);
+ req->dev.iov[2].MPID_IOV_LEN = sizeof(req->dev.stream_offset);
+ req->dev.iov_count = 3;
*buflen = sizeof(MPIDI_CH3_Pkt_t);
}
@@ -2005,7 +1999,7 @@ int MPIDI_CH3_PktPrint_Put(FILE * fp, MPIDI_CH3_Pkt_t * pkt)
MPIU_DBG_PRINTF((" addr ......... %p\n", pkt->put.addr));
MPIU_DBG_PRINTF((" count ........ %d\n", pkt->put.count));
MPIU_DBG_PRINTF((" datatype ..... 0x%08X\n", pkt->put.datatype));
- MPIU_DBG_PRINTF((" dataloop_size. 0x%08X\n", pkt->put.info.metadata.dataloop_size));
+ MPIU_DBG_PRINTF((" dataloop_size. 0x%08X\n", pkt->put.info.dataloop_size));
MPIU_DBG_PRINTF((" target ....... 0x%08X\n", pkt->put.target_win_handle));
MPIU_DBG_PRINTF((" source ....... 0x%08X\n", pkt->put.source_win_handle));
/*MPIU_DBG_PRINTF((" win_ptr ...... 0x%08X\n", pkt->put.win_ptr)); */
@@ -2018,7 +2012,7 @@ int MPIDI_CH3_PktPrint_Get(FILE * fp, MPIDI_CH3_Pkt_t * pkt)
MPIU_DBG_PRINTF((" addr ......... %p\n", pkt->get.addr));
MPIU_DBG_PRINTF((" count ........ %d\n", pkt->get.count));
MPIU_DBG_PRINTF((" datatype ..... 0x%08X\n", pkt->get.datatype));
- MPIU_DBG_PRINTF((" dataloop_size. %d\n", pkt->get.info.metadata.dataloop_size));
+ MPIU_DBG_PRINTF((" dataloop_size. %d\n", pkt->get.info.dataloop_size));
MPIU_DBG_PRINTF((" request ...... 0x%08X\n", pkt->get.request_handle));
MPIU_DBG_PRINTF((" target ....... 0x%08X\n", pkt->get.target_win_handle));
MPIU_DBG_PRINTF((" source ....... 0x%08X\n", pkt->get.source_win_handle));
@@ -2043,7 +2037,7 @@ int MPIDI_CH3_PktPrint_Accumulate(FILE * fp, MPIDI_CH3_Pkt_t * pkt)
MPIU_DBG_PRINTF((" addr ......... %p\n", pkt->accum.addr));
MPIU_DBG_PRINTF((" count ........ %d\n", pkt->accum.count));
MPIU_DBG_PRINTF((" datatype ..... 0x%08X\n", pkt->accum.datatype));
- MPIU_DBG_PRINTF((" dataloop_size. %d\n", pkt->accum.info.metadata.dataloop_size));
+ MPIU_DBG_PRINTF((" dataloop_size. %d\n", pkt->accum.info.dataloop_size));
MPIU_DBG_PRINTF((" op ........... 0x%08X\n", pkt->accum.op));
MPIU_DBG_PRINTF((" target ....... 0x%08X\n", pkt->accum.target_win_handle));
MPIU_DBG_PRINTF((" source ....... 0x%08X\n", pkt->accum.source_win_handle));
-----------------------------------------------------------------------
Summary of changes:
src/mpid/ch3/include/mpid_rma_issue.h | 26 ++++++--
src/mpid/ch3/include/mpid_rma_types.h | 3 +
src/mpid/ch3/include/mpidpkt.h | 48 +++------------
src/mpid/ch3/include/mpidrma.h | 17 +-----
src/mpid/ch3/src/ch3u_handle_recv_req.c | 82 ++++++++++++++-----------
src/mpid/ch3/src/ch3u_rma_ops.c | 26 ++++++---
src/mpid/ch3/src/ch3u_rma_pkthandler.c | 98 +++++++++++++++----------------
7 files changed, 146 insertions(+), 154 deletions(-)
hooks/post-receive
--
MPICH primary repository
1
0
[mpich] MPICH primary repository branch, master, updated. v3.2b1-82-gc09f396
by noreply@mpich.org 20 Apr '15
by noreply@mpich.org 20 Apr '15
20 Apr '15
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "MPICH primary repository".
The branch, master has been updated
via c09f396958cbef74684c3b423315b972973042a5 (commit)
from 4ef8d551fdaa38e636927cd38805d60fb13ede8d (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
http://git.mpich.org/mpich.git/commitdiff/c09f396958cbef74684c3b423315b9729…
commit c09f396958cbef74684c3b423315b972973042a5
Author: Antonio J. Pena <apenya(a)mcs.anl.gov>
Date: Mon Apr 20 14:38:26 2015 -0500
Moved a request assert to an earlier location
An assert protecting from a non-null request was happening too late in pkt_COOKIE_handler from
mpid_nem_lmt.c. This patch moves it to an earlier location so that it's checked before it's first
used.
Reported by Dmitry Polyakov.
diff --git a/src/mpid/ch3/channels/nemesis/src/mpid_nem_lmt.c b/src/mpid/ch3/channels/nemesis/src/mpid_nem_lmt.c
index c07d322..60b782d 100644
--- a/src/mpid/ch3/channels/nemesis/src/mpid_nem_lmt.c
+++ b/src/mpid/ch3/channels/nemesis/src/mpid_nem_lmt.c
@@ -462,13 +462,14 @@ static int pkt_COOKIE_handler(MPIDI_VC_t *vc, MPIDI_CH3_Pkt_t *pkt, MPIDI_msg_sz
if (cookie_pkt->from_sender) {
MPID_Request_get_ptr(cookie_pkt->receiver_req_id, req);
+ MPIU_Assert(req != NULL);
req->ch.lmt_req_id = cookie_pkt->sender_req_id;
}
else {
MPID_Request_get_ptr(cookie_pkt->sender_req_id, req);
+ MPIU_Assert(req != NULL);
req->ch.lmt_req_id = cookie_pkt->receiver_req_id;
}
- MPIU_Assert(req != NULL);
if (cookie_pkt->cookie_len != 0)
{
-----------------------------------------------------------------------
Summary of changes:
src/mpid/ch3/channels/nemesis/src/mpid_nem_lmt.c | 3 ++-
1 files changed, 2 insertions(+), 1 deletions(-)
hooks/post-receive
--
MPICH primary repository
1
0
[mpich] MPICH primary repository branch, master, updated. v3.2b1-81-g4ef8d55
by noreply@mpich.org 20 Apr '15
by noreply@mpich.org 20 Apr '15
20 Apr '15
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "MPICH primary repository".
The branch, master has been updated
via 4ef8d551fdaa38e636927cd38805d60fb13ede8d (commit)
from c89e8d8ef3f15facf9763cfba4bd7a1f30b5b16e (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
http://git.mpich.org/mpich.git/commitdiff/4ef8d551fdaa38e636927cd38805d60fb…
commit 4ef8d551fdaa38e636927cd38805d60fb13ede8d
Author: Charles J Archer <charles.j.archer(a)intel.com>
Date: Mon Apr 20 12:22:07 2015 -0700
OFI: Fix multiple providers failing probe-unexp
Update OFI netmod to match portals4 netmod anysource_matched semantics.
diff --git a/src/mpid/ch3/channels/nemesis/netmod/ofi/ofi_progress.c b/src/mpid/ch3/channels/nemesis/netmod/ofi/ofi_progress.c
index cc2048f..cdf5535 100644
--- a/src/mpid/ch3/channels/nemesis/netmod/ofi/ofi_progress.c
+++ b/src/mpid/ch3/channels/nemesis/netmod/ofi/ofi_progress.c
@@ -97,9 +97,9 @@ int MPID_nem_ofi_iprobe_impl(struct MPIDI_VC *vc,
if (rreq_ptr) {
MPIDI_CH3_Request_destroy(rreq);
*rreq_ptr = NULL;
- *flag = 0;
}
MPID_nem_ofi_poll(MPID_NONBLOCKING_POLL);
+ *flag = 0;
goto fn_exit;
}
MPIU_ERR_CHKANDJUMP4((ret < 0), mpi_errno, MPI_ERR_OTHER,
diff --git a/src/mpid/ch3/channels/nemesis/netmod/ofi/ofi_tagged.c b/src/mpid/ch3/channels/nemesis/netmod/ofi/ofi_tagged.c
index c73f618..66b9130 100644
--- a/src/mpid/ch3/channels/nemesis/netmod/ofi/ofi_tagged.c
+++ b/src/mpid/ch3/channels/nemesis/netmod/ofi/ofi_tagged.c
@@ -68,10 +68,9 @@ static inline int MPID_nem_ofi_recv_callback(cq_tagged_entry_t * wc, MPID_Reques
/* ---------------------------------------------------- */
rreq->status.MPI_ERROR = MPI_SUCCESS;
rreq->status.MPI_SOURCE = get_source(wc->tag);
- rreq->status.MPI_TAG = get_tag(wc->tag);
+ rreq->status.MPI_TAG = get_tag(wc->tag);
REQ_OFI(rreq)->req_started = 1;
MPIR_STATUS_SET_COUNT(rreq->status, wc->len);
-
if (REQ_OFI(rreq)->pack_buffer) {
MPIDI_CH3U_Buffer_copy(REQ_OFI(rreq)->pack_buffer,
MPIR_STATUS_GET_COUNT(rreq->status),
@@ -388,7 +387,7 @@ void MPID_nem_ofi_anysource_posted(MPID_Request * rreq)
#define FCNAME DECL_FUNC(MPID_nem_ofi_anysource_matched)
int MPID_nem_ofi_anysource_matched(MPID_Request * rreq)
{
- int mpi_errno = FALSE;
+ int mpi_errno = TRUE;
int ret;
BEGIN_FUNC(FCNAME);
/* ----------------------------------------------------- */
@@ -401,10 +400,7 @@ int MPID_nem_ofi_anysource_matched(MPID_Request * rreq)
/* --------------------------------------------------- */
/* Request cancelled: cancel and complete the request */
/* --------------------------------------------------- */
- mpi_errno = TRUE;
- MPIR_STATUS_SET_CANCEL_BIT(rreq->status, TRUE);
- MPIR_STATUS_SET_COUNT(rreq->status, 0);
- MPIDI_CH3U_Request_complete(rreq);
+ mpi_errno = FALSE;
}
END_FUNC(FCNAME);
return mpi_errno;
-----------------------------------------------------------------------
Summary of changes:
.../ch3/channels/nemesis/netmod/ofi/ofi_progress.c | 2 +-
.../ch3/channels/nemesis/netmod/ofi/ofi_tagged.c | 10 +++-------
2 files changed, 4 insertions(+), 8 deletions(-)
hooks/post-receive
--
MPICH primary repository
1
0
[mpich] MPICH primary repository branch, master, updated. v3.2b1-80-gc89e8d8
by noreply@mpich.org 20 Apr '15
by noreply@mpich.org 20 Apr '15
20 Apr '15
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "MPICH primary repository".
The branch, master has been updated
via c89e8d8ef3f15facf9763cfba4bd7a1f30b5b16e (commit)
from cae2623451508fa68f2167f00debdb218c7989aa (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
http://git.mpich.org/mpich.git/commitdiff/c89e8d8ef3f15facf9763cfba4bd7a1f3…
commit c89e8d8ef3f15facf9763cfba4bd7a1f30b5b16e
Author: Charles J Archer <charles.j.archer(a)intel.com>
Date: Mon Apr 20 08:25:36 2015 -0700
OFI: Update for removed FI_CANCEL flag
diff --git a/src/mpid/ch3/channels/nemesis/netmod/ofi/ofi_init.c b/src/mpid/ch3/channels/nemesis/netmod/ofi/ofi_init.c
index 3993088..91576f6 100644
--- a/src/mpid/ch3/channels/nemesis/netmod/ofi/ofi_init.c
+++ b/src/mpid/ch3/channels/nemesis/netmod/ofi/ofi_init.c
@@ -81,7 +81,6 @@ int MPID_nem_ofi_init(MPIDI_PG_t * pg_p, int pg_rank, char **bc_val_p, int *val_
hints->mode = FI_CONTEXT;
hints->ep_attr->type = FI_EP_RDM; /* Reliable datagram */
hints->caps = FI_TAGGED; /* Tag matching interface */
- hints->caps |= FI_CANCEL; /* Support cancel */
hints->caps |= FI_DYNAMIC_MR; /* Global dynamic mem region */
/* ------------------------------------------------------------------------ */
-----------------------------------------------------------------------
Summary of changes:
.../ch3/channels/nemesis/netmod/ofi/ofi_init.c | 1 -
1 files changed, 0 insertions(+), 1 deletions(-)
hooks/post-receive
--
MPICH primary repository
1
0
[mpich] MPICH primary repository branch, master, updated. v3.2b1-79-gcae2623
by noreply@mpich.org 17 Apr '15
by noreply@mpich.org 17 Apr '15
17 Apr '15
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "MPICH primary repository".
The branch, master has been updated
via cae2623451508fa68f2167f00debdb218c7989aa (commit)
from 17e31e59b491fc1383653ddd41f8264e8c5c69e6 (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
http://git.mpich.org/mpich.git/commitdiff/cae2623451508fa68f2167f00debdb218…
commit cae2623451508fa68f2167f00debdb218c7989aa
Author: Halim Amer <aamer(a)anl.gov>
Date: Fri Apr 17 19:20:32 2015 -0500
Revert "Applied the PPoPP patch"
This reverts commit 17e31e59b491fc1383653ddd41f8264e8c5c69e6.
diff --git a/src/include/mpiimplthreadpost.h b/src/include/mpiimplthreadpost.h
index 60a0882..be4578f 100644
--- a/src/include/mpiimplthreadpost.h
+++ b/src/include/mpiimplthreadpost.h
@@ -191,13 +191,7 @@ MPIU_Thread_CS_yield_lockname_recursive_impl_(enum MPIU_Nest_mutexes kind,
MPID_Thread_mutex_unlock(mutex);
MPID_Thread_yield();
-
-/* Use low priority here because this thread has a lower probability
- * to do useful work compared to others outside the progress engine.
- * This is only effective with the priority lock.
- * The other lock types implement a plain lock underneath.
- */
- MPID_Thread_mutex_lock_low(mutex);
+ MPID_Thread_mutex_lock(mutex);
}
/* undef for safety, this is a commonly-included header */
diff --git a/src/include/thread/mpiu_thread_posix_funcs.h b/src/include/thread/mpiu_thread_posix_funcs.h
index facea60..10af9fc 100644
--- a/src/include/thread/mpiu_thread_posix_funcs.h
+++ b/src/include/thread/mpiu_thread_posix_funcs.h
@@ -8,10 +8,8 @@
/*
* Threads
*/
-#include <limits.h>
+
#include "mpiu_process_wrappers.h" /* for MPIU_PW_Sched_yield */
-#include "mpidbg.h"
-#include "opa_primitives.h"
/*
One of PTHREAD_MUTEX_RECURSIVE_NP and PTHREAD_MUTEX_RECURSIVE seem to be
@@ -55,443 +53,144 @@ do { \
MPIU_DBG_MSG(THREAD,VERBOSE,"exit MPIU_Thread_yield"); \
} while (0)
-/*----------------------*/
-/* Ticket Lock Routines */
-/*----------------------*/
-
-static inline int ticket_lock_init(ticket_lock_t *lock)
-{
- OPA_store_int(&lock->next_ticket, 0);
- OPA_store_int(&lock->now_serving, 0);
- return 0;
-}
-
-/* Atomically increment the nex_ticket counter and get my ticket.
- Then spin on now_serving until it equals my ticket.
- */
-static inline int ticket_acquire_lock(ticket_lock_t* lock)
-{
- int my_ticket = OPA_fetch_and_add_int(&lock->next_ticket, 1);
- while(OPA_load_int(&lock->now_serving) != my_ticket)
- ;
- return 0;
-}
-
-/* Release the lock
- */
-static inline int ticket_release_lock(ticket_lock_t* lock)
-{
- /* Avoid compiler reordering before releasing the lock*/
- OPA_compiler_barrier();
- OPA_incr_int(&lock->now_serving);
- return 0;
-}
-
-/*------------------------*/
-/* Priority Lock Routines */
-/*------------------------*/
-
-static inline int priority_lock_init(priority_lock_t *lock)
-{
- OPA_store_int(&lock->next_ticket_H, 0);
- OPA_store_int(&lock->now_serving_H, 0);
- OPA_store_int(&lock->next_ticket_L, 0);
- OPA_store_int(&lock->now_serving_L, 0);
- OPA_store_int(&lock->next_ticket_B, 0);
- OPA_store_int(&lock->now_serving_B, 0);
- lock->already_blocked = 0;
- return 0;
-}
-
-/* First wait my turn in this priority level, Then if I am the first
- one to block the LPRs, wait for the last LPR to terminate.
- */
-static inline int priority_acquire_lock(priority_lock_t* lock)
-{
- int my_ticket = OPA_fetch_and_add_int(&lock->next_ticket_H, 1);
- while(OPA_load_int(&lock->now_serving_H) != my_ticket)
- ;
- int B_ticket;
- if(!lock->already_blocked)
- {
- B_ticket = OPA_fetch_and_add_int(&lock->next_ticket_B, 1);
- while(OPA_load_int(&lock->now_serving_B) != B_ticket)
- ;
- lock->already_blocked = 1;
- }
- lock->last_acquisition_priority = HIGH_PRIORITY;
- return 0;
-}
-static inline int priority_acquire_lock_low(priority_lock_t* lock)
-{
- int my_ticket = OPA_fetch_and_add_int(&lock->next_ticket_L, 1);
- while(OPA_load_int(&lock->now_serving_L) != my_ticket)
- ;
- int B_ticket = OPA_fetch_and_add_int(&lock->next_ticket_B, 1);
- while(OPA_load_int(&lock->now_serving_B) != B_ticket)
- ;
- lock->last_acquisition_priority = LOW_PRIORITY;
- return 0;
-}
-
-/* Release the lock
- */
-static inline int priority_release_lock(priority_lock_t* lock)
-{
- int err=0;
- if(lock->last_acquisition_priority==HIGH_PRIORITY)
- err = priority_release_lock_high(lock);
- else
- err = priority_release_lock_low(lock);
- return err;
-}
-
-static inline int priority_release_lock_high(priority_lock_t* lock)
-{
- /* Avoid compiler reordering before releasing the lock*/
- OPA_compiler_barrier();
- /* Only me in the HPRs queue -> let the LPRs pass */
- if(OPA_load_int(&lock->now_serving_H) == OPA_load_int(&lock->next_ticket_H) - 1)
- {
- lock->already_blocked = 0;
- OPA_incr_int(&lock->now_serving_B);
- }
- OPA_incr_int(&lock->now_serving_H);
- return 0;
-}
-
-static inline int priority_release_lock_low(priority_lock_t* lock)
-{
- /* Avoid compiler reordering before releasing the lock*/
- OPA_compiler_barrier();
- OPA_incr_int(&lock->now_serving_B);
- OPA_incr_int(&lock->now_serving_L);
- return 0;
-}
/*
- * MPIU Mutexes: encapsulate lower level locks like pthread mutexes
- * and ticket locks
+ * Mutexes
+ */
+
+/* FIXME: mutex creation and destruction should be implemented as routines
+ because there is no reason to use macros (these are not on the performance
+ critical path). Making these macros requires that any code that might use
+ these must load all of the pthread.h (or other thread library) support.
*/
+/* FIXME: using constant initializer if available */
#if !defined(MPICH_DEBUG_MUTEX) || !defined(PTHREAD_MUTEX_ERRORCHECK_VALUE)
-static inline void MPIU_Thread_mutex_create(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
-{
- int err__=0;
- switch(MPIU_lock_type){
- case MPIU_MUTEX:
- {
- err__ = pthread_mutex_init(&mutex_ptr_->pthread_lock, NULL);
- break;
- }
- case MPIU_TICKET:
- {
- err__ = ticket_lock_init(&mutex_ptr_->ticket_lock);
- break;
- }
- case MPIU_PRIORITY:
- {
- err__ = priority_lock_init(&mutex_ptr_->priority_lock);
- break;
- }
- }
- /* FIXME: convert error to an MPIU_THREAD_ERR value */
- *(int *)(err_ptr_) = err__;
- MPIU_DBG_MSG_P(THREAD,TYPICAL,"Created MPIU_Thread_mutex %p", (mutex_ptr_));
-}
+#define MPIU_Thread_mutex_create(mutex_ptr_, err_ptr_) \
+do { \
+ int err__; \
+ \
+ err__ = pthread_mutex_init((mutex_ptr_), NULL); \
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
+ *(int *)(err_ptr_) = err__; \
+ MPIU_DBG_MSG_P(THREAD,TYPICAL,"Created MPIU_Thread_mutex %p", (mutex_ptr_)); \
+} while (0)
#else /* MPICH_DEBUG_MUTEX */
-static inline void MPIU_Thread_mutex_create(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
-{
- int err__=0;
- switch(MPIU_lock_type){
- case MPIU_MUTEX:
- {
- pthread_mutexattr_t attr__;
- /* FIXME this used to be PTHREAD_MUTEX_ERRORCHECK_NP, but we had to change
- it for the thread granularity work when we needed recursive mutexes. We
- should go through this code and see if there's any good way to implement
- error checked versions with the recursive mutexes. */
- pthread_mutexattr_init(&attr__);
- pthread_mutexattr_settype(&attr__, PTHREAD_MUTEX_ERRORCHECK_VALUE);
- err__ = pthread_mutex_init(&mutex_ptr_->pthread_lock, &attr__);
- break;
- }
- case MPIU_TICKET:
- {
- err__ = ticket_lock_init(&mutex_ptr_->ticket_lock);
- break;
- }
- case MPIU_PRIORITY:
- {
- err__ = priority_lock_init(&mutex_ptr_->priority_lock);
- break;
- }
- }
- if (err__)
+#define MPIU_Thread_mutex_create(mutex_ptr_, err_ptr_) \
+do { \
+ int err__; \
+ pthread_mutexattr_t attr__; \
+ \
+ /* FIXME this used to be PTHREAD_MUTEX_ERRORCHECK_NP, but we had to change
+ it for the thread granularity work when we needed recursive mutexes. We
+ should go through this code and see if there's any good way to implement
+ error checked versions with the recursive mutexes. */ \
+ pthread_mutexattr_init(&attr__); \
+ pthread_mutexattr_settype(&attr__, PTHREAD_MUTEX_ERRORCHECK_VALUE); \
+ err__ = pthread_mutex_init((mutex_ptr_), &attr__); \
+ if (err__) \
MPIU_Internal_sys_error_printf("pthread_mutex_init", err__, \
- " %s:%d\n", __FILE__, __LINE__);
- /* FIXME: convert error to an MPIU_THREAD_ERR value */
- *(int *)(err_ptr_) = err__;
- MPIU_DBG_MSG_P(THREAD,TYPICAL,"Created MPIU_Thread_mutex %p", (mutex_ptr_));
-}
+ " %s:%d\n", __FILE__, __LINE__);\
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
+ *(int *)(err_ptr_) = err__; \
+ MPIU_DBG_MSG_P(THREAD,TYPICAL,"Created MPIU_Thread_mutex %p", (mutex_ptr_)); \
+} while (0)
#endif
-static inline void MPIU_Thread_mutex_destroy(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
-{
- int err__=0;
-
- MPIU_DBG_MSG_P(THREAD,TYPICAL,"About to destroy MPIU_Thread_mutex %p", (mutex_ptr_));
-
- switch(MPIU_lock_type){
- case MPIU_MUTEX:
- {
- err__ = pthread_mutex_destroy(&mutex_ptr_->pthread_lock);
- break;
- }
- case MPIU_TICKET:
- {
- /* FIXME Should we do something here ?*/
- err__ = 0;
- break;
- }
- case MPIU_PRIORITY:
- {
- /* FIXME Should we do something here ?*/
- err__ = 0;
- break;
- }
- }
- /* FIXME: convert error to an MPIU_THREAD_ERR value */
- *(int *)(err_ptr_) = err__;
-}
-
-#ifndef MPICH_DEBUG_MUTEX
-static inline void MPIU_Thread_mutex_lock(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
-{
- int err__=0;
- MPIU_DBG_MSG_P(THREAD,VERBOSE,"enter MPIU_Thread_mutex_lock %p", (mutex_ptr_));
-
- switch(MPIU_lock_type){
- case MPIU_MUTEX:
- {
- err__ = pthread_mutex_lock(&mutex_ptr_->pthread_lock);
- break;
- }
- case MPIU_TICKET:
- {
- err__ = ticket_acquire_lock(&mutex_ptr_->ticket_lock);
- break;
- }
- case MPIU_PRIORITY:
- {
- err__ = priority_acquire_lock(&mutex_ptr_->priority_lock);
- break;
- }
- }
-
- /* FIXME: convert error to an MPIU_THREAD_ERR value */
- *(int *)(err_ptr_) = err__;
- MPIU_DBG_MSG_P(THREAD,VERBOSE,"exit MPIU_Thread_mutex_lock %p", (mutex_ptr_));
-}
-#else /* MPICH_DEBUG_MUTEX */
-static inline void MPIU_Thread_mutex_lock(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
-{
- int err__=0;
- MPIU_DBG_MSG_P(THREAD,VERBOSE,"enter MPIU_Thread_mutex_lock %p", (mutex_ptr_));
- switch(MPIU_lock_type){
- case MPIU_MUTEX:
- {
- err__ = pthread_mutex_lock(&mutex_ptr_->pthread_lock);
- break;
- }
- case MPIU_TICKET:
- {
- err__ = ticket_acquire_lock(&mutex_ptr_->ticket_lock);
- break;
- }
- case MPIU_PRIORITY:
- {
- err__ = priority_acquire_lock(&mutex_ptr_->priority_lock);
- break;
- }
- }
-
- if (err__)
- {
- MPIU_DBG_MSG_S(THREAD,TERSE," mutex lock error: %s", MPIU_Strerror(err__));
- MPIU_Internal_sys_error_printf("pthread_mutex_lock", err__,\
- " %s:%d\n", __FILE__, __LINE__);
- }
- /* FIXME: convert error to an MPIU_THREAD_ERR value */
- *(int *)(err_ptr_) = err__;
- MPIU_DBG_MSG_P(THREAD,VERBOSE,"exit MPIU_Thread_mutex_lock %p", (mutex_ptr_));
-}
-#endif
+#define MPIU_Thread_mutex_destroy(mutex_ptr_, err_ptr_) \
+do { \
+ int err__; \
+ \
+ MPIU_DBG_MSG_P(THREAD,TYPICAL,"About to destroy MPIU_Thread_mutex %p", (mutex_ptr_)); \
+ err__ = pthread_mutex_destroy(mutex_ptr_); \
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
+ *(int *)(err_ptr_) = err__; \
+} while (0)
#ifndef MPICH_DEBUG_MUTEX
-static inline void MPIU_Thread_mutex_lock_low(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
-{
- int err__=0;
- MPIU_DBG_MSG_P(THREAD,VERBOSE,"enter MPIU_Thread_mutex_lock %p", (mutex_ptr_));
-
- switch(MPIU_lock_type){
- case MPIU_MUTEX:
- {
- err__ = pthread_mutex_lock(&mutex_ptr_->pthread_lock);
- break;
- }
- case MPIU_TICKET:
- {
- err__ = ticket_acquire_lock(&mutex_ptr_->ticket_lock);
- break;
- }
- case MPIU_PRIORITY:
- {
- err__ = priority_acquire_lock_low(&mutex_ptr_->priority_lock);
- break;
- }
- }
-
- /* FIXME: convert error to an MPIU_THREAD_ERR value */
- *(int *)(err_ptr_) = err__;
- MPIU_DBG_MSG_P(THREAD,VERBOSE,"exit MPIU_Thread_mutex_lock %p", (mutex_ptr_));
-}
+#define MPIU_Thread_mutex_lock(mutex_ptr_, err_ptr_) \
+do { \
+ int err__; \
+ MPIU_DBG_MSG_P(THREAD,VERBOSE,"enter MPIU_Thread_mutex_lock %p", (mutex_ptr_)); \
+ err__ = pthread_mutex_lock(mutex_ptr_); \
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
+ *(int *)(err_ptr_) = err__; \
+ MPIU_DBG_MSG_P(THREAD,VERBOSE,"exit MPIU_Thread_mutex_lock %p", (mutex_ptr_)); \
+} while (0)
#else /* MPICH_DEBUG_MUTEX */
-static inline void MPIU_Thread_mutex_lock_low(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
-{
- int err__=0;
- MPIU_DBG_MSG_P(THREAD,VERBOSE,"enter MPIU_Thread_mutex_lock %p", (mutex_ptr_));
- switch(MPIU_lock_type){
- case MPIU_MUTEX:
- {
- err__ = pthread_mutex_lock(&mutex_ptr_->pthread_lock);
- break;
- }
- case MPIU_TICKET:
- {
- err__ = ticket_acquire_lock(&mutex_ptr_->ticket_lock);
- break;
- }
- case MPIU_PRIORITY:
- {
- err__ = priority_acquire_lock_low(&mutex_ptr_->priority_lock);
- break;
- }
- }
-
- if (err__)
- {
- MPIU_DBG_MSG_S(THREAD,TERSE," mutex lock error: %s", MPIU_Strerror(err__));
+#define MPIU_Thread_mutex_lock(mutex_ptr_, err_ptr_) \
+do { \
+ int err__; \
+ MPIU_DBG_MSG_P(THREAD,VERBOSE,"enter MPIU_Thread_mutex_lock %p", (mutex_ptr_)); \
+ err__ = pthread_mutex_lock(mutex_ptr_); \
+ if (err__) \
+ { \
+ MPIU_DBG_MSG_S(THREAD,TERSE," mutex lock error: %s", MPIU_Strerror(err__)); \
MPIU_Internal_sys_error_printf("pthread_mutex_lock", err__,\
- " %s:%d\n", __FILE__, __LINE__);
- }
- /* FIXME: convert error to an MPIU_THREAD_ERR value */
- *(int *)(err_ptr_) = err__;
- MPIU_DBG_MSG_P(THREAD,VERBOSE,"exit MPIU_Thread_mutex_lock %p", (mutex_ptr_));
-}
+ " %s:%d\n", __FILE__, __LINE__);\
+ } \
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
+ *(int *)(err_ptr_) = err__; \
+ MPIU_DBG_MSG_P(THREAD,VERBOSE,"exit MPIU_Thread_mutex_lock %p", (mutex_ptr_)); \
+} while (0)
#endif
#ifndef MPICH_DEBUG_MUTEX
-static inline void MPIU_Thread_mutex_unlock(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
-{
- int err__=0;
-
- MPIU_DBG_MSG_P(THREAD,TYPICAL,"MPIU_Thread_mutex_unlock %p", (mutex_ptr_));
-
- switch(MPIU_lock_type){
- case MPIU_MUTEX:
- {
- err__ = pthread_mutex_unlock(&mutex_ptr_->pthread_lock);
- break;
- }
- case MPIU_TICKET:
- {
- err__ = ticket_release_lock(&mutex_ptr_->ticket_lock);
- break;
- }
- case MPIU_PRIORITY:
- {
- err__ = priority_release_lock(&mutex_ptr_->priority_lock);
- break;
- }
- }
-
- /* FIXME: convert error to an MPIU_THREAD_ERR value */
- *(int *)(err_ptr_) = err__;
-}
+#define MPIU_Thread_mutex_unlock(mutex_ptr_, err_ptr_) \
+do { \
+ int err__; \
+ \
+ MPIU_DBG_MSG_P(THREAD,TYPICAL,"MPIU_Thread_mutex_unlock %p", (mutex_ptr_)); \
+ err__ = pthread_mutex_unlock(mutex_ptr_); \
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
+ *(int *)(err_ptr_) = err__; \
+} while (0)
#else /* MPICH_DEBUG_MUTEX */
-static inline void MPIU_Thread_mutex_unlock(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
-{
- int err__=0;
-
- MPIU_DBG_MSG_P(THREAD,VERBOSE,"MPIU_Thread_mutex_unlock %p", (mutex_ptr_));
-
- switch(MPIU_lock_type){
- case MPIU_MUTEX:
- {
- err__ = pthread_mutex_unlock(&mutex_ptr_->pthread_lock);
- break;
- }
- case MPIU_TICKET:
- {
- err__ = ticket_release_lock(&mutex_ptr_->ticket_lock);
- break;
- }
- case MPIU_PRIORITY:
- {
- err__ = priority_release_lock(&mutex_ptr_->priority_lock);
- break;
- }
- }
- if (err__)
- {
- MPIU_DBG_MSG_S(THREAD,TERSE," mutex unlock error: %s", MPIU_Strerror(err__));
- MPIU_Internal_sys_error_printf("pthread_mutex_unlock", err__,\
- " %s:%d\n", __FILE__, __LINE__);
- }
- /* FIXME: convert error to an MPIU_THREAD_ERR value */
- *(int *)(err_ptr_) = err__;
-}
+#define MPIU_Thread_mutex_unlock(mutex_ptr_, err_ptr_) \
+do { \
+ int err__; \
+ \
+ MPIU_DBG_MSG_P(THREAD,VERBOSE,"MPIU_Thread_mutex_unlock %p", (mutex_ptr_)); \
+ err__ = pthread_mutex_unlock(mutex_ptr_); \
+ if (err__) \
+ { \
+ MPIU_DBG_MSG_S(THREAD,TERSE," mutex unlock error: %s", MPIU_Strerror(err__)); \
+ MPIU_Internal_sys_error_printf("pthread_mutex_unlock", err__, \
+ " %s:%d\n", __FILE__, __LINE__); \
+ } \
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
+ *(int *)(err_ptr_) = err__; \
+} while (0)
#endif
#ifndef MPICH_DEBUG_MUTEX
-static inline void MPIU_Thread_mutex_trylock(MPIU_Thread_mutex_t* mutex_ptr_, int* flag_ptr_, int* err_ptr_)
-{
- int err__=0;
-
- if(MPIU_lock_type == MPIU_MUTEX)
- err__ = pthread_mutex_trylock(&mutex_ptr_->pthread_lock);
- else
- /*No trylock routine for ticket-based locks*/
- err__ = 1;
-
- *(flag_ptr_) = (err__ == 0) ? TRUE : FALSE;
- MPIU_DBG_MSG_FMT(THREAD,VERBOSE,(MPIU_DBG_FDEST, "MPIU_Thread_mutex_trylock mutex=%p result=%s", (mutex_ptr_), (*(flag_ptr_) ? "success" : "failure")));
- *(int *)(err_ptr_) = (err__ == EBUSY) ? MPIU_THREAD_SUCCESS : err__;
- /* FIXME: convert error to an MPIU_THREAD_ERR value */
-}
+#define MPIU_Thread_mutex_trylock(mutex_ptr_, flag_ptr_, err_ptr_) \
+do { \
+ int err__; \
+ \
+ err__ = pthread_mutex_trylock(mutex_ptr_); \
+ *(flag_ptr_) = (err__ == 0) ? TRUE : FALSE; \
+ MPIU_DBG_MSG_FMT(THREAD,VERBOSE,(MPIU_DBG_FDEST, "MPIU_Thread_mutex_trylock mutex=%p result=%s", (mutex_ptr_), (*(flag_ptr_) ? "success" : "failure"))); \
+ *(int *)(err_ptr_) = (err__ == EBUSY) ? MPIU_THREAD_SUCCESS : err__; \
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
+} while (0)
#else /* MPICH_DEBUG_MUTEX */
-static inline void MPIU_Thread_mutex_trylock(MPIU_Thread_mutex_t* mutex_ptr_, int* flag_ptr_, int* err_ptr_)
-{
- int err__=0;
-
- if(MPIU_lock_type == MPIU_MUTEX)
- err__ = pthread_mutex_trylock(&mutex_ptr_->pthread_lock);
- else
- /*No trylock routine for ticket-based locks*/
- err__ = 1;
-
- if (err__ && err__ != EBUSY)
- {
- MPIU_DBG_MSG_S(THREAD,TERSE," mutex trylock error: %s", MPIU_Strerror(err__));
+#define MPIU_Thread_mutex_trylock(mutex_ptr_, flag_ptr_, err_ptr_) \
+do { \
+ int err__; \
+ \
+ err__ = pthread_mutex_trylock(mutex_ptr_); \
+ if (err__ && err__ != EBUSY) \
+ { \
+ MPIU_DBG_MSG_S(THREAD,TERSE," mutex trylock error: %s", MPIU_Strerror(err__)); \
MPIU_Internal_sys_error_printf("pthread_mutex_trylock", err__,\
- " %s:%d\n", __FILE__, __LINE__);
- }
- *(flag_ptr_) = (err__ == 0) ? TRUE : FALSE;
- MPIU_DBG_MSG_FMT(THREAD,VERBOSE,(MPIU_DBG_FDEST, "MPIU_Thread_mutex_trylock mutex=%p result=%s", (mutex_ptr_), (*(flag_ptr_) ? "success" : "failure")));
- *(int *)(err_ptr_) = (err__ == EBUSY) ? MPIU_THREAD_SUCCESS : err__;
- /* FIXME: convert error to an MPIU_THREAD_ERR value */
-}
+ " %s:%d\n", __FILE__, __LINE__);\
+ } \
+ *(flag_ptr_) = (err__ == 0) ? TRUE : FALSE; \
+ MPIU_DBG_MSG_FMT(THREAD,VERBOSE,(MPIU_DBG_FDEST, "MPIU_Thread_mutex_trylock mutex=%p result=%s", (mutex_ptr_), (*(flag_ptr_) ? "success" : "failure"))); \
+ *(int *)(err_ptr_) = (err__ == EBUSY) ? MPIU_THREAD_SUCCESS : err__; \
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
+} while (0)
#endif
/*
diff --git a/src/include/thread/mpiu_thread_posix_types.h b/src/include/thread/mpiu_thread_posix_types.h
index c163b67..b38f7a1 100644
--- a/src/include/thread/mpiu_thread_posix_types.h
+++ b/src/include/thread/mpiu_thread_posix_types.h
@@ -7,59 +7,8 @@
#include <errno.h>
#include <pthread.h>
-#include "opa_primitives.h"
-
-/* Define lock types here */
-typedef enum{
- MPIU_MUTEX,
- MPIU_TICKET,
- MPIU_PRIORITY
-}MPIU_Thread_lock_impl_t;
-
-MPIU_Thread_lock_impl_t MPIU_lock_type;
-
-/*----------------------------*/
-/* Ticket lock data structure */
-/*----------------------------*/
-
-/* Define the lock as a structure of two counters */
-typedef struct ticket_lock_t{
- OPA_int_t next_ticket;
- OPA_int_t now_serving;
-} ticket_lock_t;
-
-/*------------------------------*/
-/* Priority lock data structure */
-/*------------------------------*/
-
-#define HIGH_PRIORITY 1
-#define LOW_PRIORITY 2
-
- /* We only define for now 2 levels of priority: */
- /* The high priority requests (HPRs) */
- /* The low priority requests (LPRs) */
-typedef struct priority_lock_t{
- OPA_int_t next_ticket_H __attribute__((aligned(64)));
- OPA_int_t now_serving_H;
- OPA_int_t next_ticket_L __attribute__((aligned(64)));
- OPA_int_t now_serving_L;
- /* In addition we include two other counters */
- /* so high priority requests block the lower ones */
- OPA_int_t next_ticket_B __attribute__((aligned(64)));
- OPA_int_t now_serving_B;
- /* This is to allow high priority requests know */
- /* that low priority requests are already blocked */
- /* by another high*/
- unsigned already_blocked;
- unsigned last_acquisition_priority;
-} priority_lock_t;
-
-typedef struct MPIU_Thread_mutex_t{
- pthread_mutex_t pthread_lock;
- ticket_lock_t ticket_lock;
- priority_lock_t priority_lock;
-} MPIU_Thread_mutex_t;
+typedef pthread_mutex_t MPIU_Thread_mutex_t;
typedef pthread_cond_t MPIU_Thread_cond_t;
typedef pthread_t MPIU_Thread_id_t;
typedef pthread_key_t MPIU_Thread_tls_t;
diff --git a/src/mpi/init/async.c b/src/mpi/init/async.c
index a95e648..c003827 100644
--- a/src/mpi/init/async.c
+++ b/src/mpi/init/async.c
@@ -61,10 +61,9 @@ static void progress_fn(void * data)
MPIU_Thread_mutex_unlock(&progress_mutex, &mpi_errno);
MPIU_Assert(!mpi_errno);
- if(MPIU_lock_type==MPIU_MUTEX){
- MPIU_Thread_cond_signal(&progress_cond, &mpi_errno);
- MPIU_Assert(!mpi_errno);
- } /* Else: busy loop will automatically break*/
+
+ MPIU_Thread_cond_signal(&progress_cond, &mpi_errno);
+ MPIU_Assert(!mpi_errno);
MPIU_THREAD_CS_EXIT(ALLFUNC,);
@@ -138,21 +137,16 @@ int MPIR_Finalize_async_thread(void)
/* XXX DJG why is this unlock/lock necessary? Should we just YIELD here or later? */
MPIU_THREAD_CS_EXIT(ALLFUNC,);
- if(MPIU_lock_type==MPIU_MUTEX){
- MPIU_Thread_mutex_lock(&progress_mutex, &mpi_errno);
- MPIU_Assert(!mpi_errno);
-
- while (!progress_thread_done) {
- MPIU_Thread_cond_wait(&progress_cond, &progress_mutex.pthread_lock, &mpi_errno);
- MPIU_Assert(!mpi_errno);
- }
-
- MPIU_Thread_mutex_unlock(&progress_mutex, &mpi_errno);
- MPIU_Assert(!mpi_errno);
- }
- else
- while (!progress_thread_done) ; /* busy loop */
- /* No need to unlock the mutex */
+ MPIU_Thread_mutex_lock(&progress_mutex, &mpi_errno);
+ MPIU_Assert(!mpi_errno);
+
+ while (!progress_thread_done) {
+ MPIU_Thread_cond_wait(&progress_cond, &progress_mutex, &mpi_errno);
+ MPIU_Assert(!mpi_errno);
+ }
+
+ MPIU_Thread_mutex_unlock(&progress_mutex, &mpi_errno);
+ MPIU_Assert(!mpi_errno);
mpi_errno = MPIR_Comm_free_impl(progress_comm_ptr);
MPIU_Assert(!mpi_errno);
diff --git a/src/mpi/init/initthread.c b/src/mpi/init/initthread.c
index 2b5db44..aaf20c8 100644
--- a/src/mpi/init/initthread.c
+++ b/src/mpi/init/initthread.c
@@ -190,16 +190,6 @@ static int MPIR_Thread_CS_Init( void )
int err;
MPIU_THREADPRIV_DECL;
- MPIU_lock_type = MPIU_MUTEX;
- char *s;
- s = getenv( "MPICH_LOCK_TYPE" );
- if(s){
- if(strcmp( "ticket", s ) == 0)
- MPIU_lock_type = MPIU_TICKET;
- else if (strcmp( "priority", s ) == 0)
- MPIU_lock_type = MPIU_PRIORITY;
- }
-
MPIU_Assert(MPICH_MAX_LOCKS >= MPIU_Nest_NUM_MUTEXES);
/* we create this at all granularities right now */
@@ -568,30 +558,6 @@ int MPIR_Init_thread(int * argc, char ***argv, int required, int * provided)
if (mpi_errno == MPI_SUCCESS)
mpi_errno = MPID_InitCompleted();
- const char *s;
-
- switch(MPIU_lock_type){
- case MPIU_MUTEX:
- {
- s = "mutex";
- break;
- }
- case MPIU_TICKET:
- {
- s = "ticket";
- break;
- }
- case MPIU_PRIORITY:
- {
- s = "priority";
- break;
- }
- }
-
- if(MPIR_Process.comm_world->rank==0){
- printf("\n[MPICH INFO] Critical section(s) based on %s \n\n", s);
- }
-
fn_exit:
MPIU_THREAD_CS_EXIT(INIT,required);
/* Make fields of MPIR_Process global visible and set mpich_state
@@ -660,7 +626,6 @@ int MPI_Init_thread( int *argc, char ***argv, int required, int *provided )
{
int mpi_errno = MPI_SUCCESS;
int rc ATTRIBUTE((unused)), reqd = required;
-
MPID_MPI_INIT_STATE_DECL(MPID_STATE_MPI_INIT_THREAD);
rc = MPID_Wtime_init();
diff --git a/src/mpid/ch3/channels/nemesis/src/ch3_progress.c b/src/mpid/ch3/channels/nemesis/src/ch3_progress.c
index 8dda2cd..f022651 100644
--- a/src/mpid/ch3/channels/nemesis/src/ch3_progress.c
+++ b/src/mpid/ch3/channels/nemesis/src/ch3_progress.c
@@ -563,24 +563,13 @@ static int MPIDI_CH3I_Progress_delay(unsigned int completion_count)
/* FIXME should be appropriately abstracted somehow */
# if defined(MPICH_IS_THREADED) && (MPIU_THREAD_GRANULARITY == MPIU_THREAD_GRANULARITY_GLOBAL)
{
- if(MPIU_lock_type==MPIU_MUTEX){
while (1)
{
if (completion_count != OPA_load_int(&MPIDI_CH3I_progress_completion_count) ||
MPIDI_CH3I_progress_blocked != TRUE)
break;
- MPID_Thread_cond_wait(&MPIDI_CH3I_progress_completion_cond, &MPIR_ThreadInfo.global_mutex.pthread_lock/*MPIDCOMM*/);
+ MPID_Thread_cond_wait(&MPIDI_CH3I_progress_completion_cond, &MPIR_ThreadInfo.global_mutex/*MPIDCOMM*/);
}
- }
- else{
- /* First release the lock and then enter a busy loop*/
- MPID_Thread_mutex_unlock(&MPIR_ThreadInfo.global_mutex/*MPIDCOMM*/);
- while (completion_count == OPA_load_int(&MPIDI_CH3I_progress_completion_count) &&
- MPIDI_CH3I_progress_blocked == TRUE)
- ; /* wait in a busy loop*/
- /* Hold the lock again*/
- MPID_Thread_mutex_lock(&MPIR_ThreadInfo.global_mutex/*MPIDCOMM*/);
- }
}
# endif
@@ -604,9 +593,7 @@ static int MPIDI_CH3I_Progress_continue(unsigned int completion_count/*unused*/)
# if defined(MPICH_IS_THREADED) && (MPIU_THREAD_GRANULARITY == MPIU_THREAD_GRANULARITY_GLOBAL)
{
/* we currently hold the MPIDCOMM CS */
- if(MPIU_lock_type==MPIU_MUTEX)
MPID_Thread_cond_broadcast(&MPIDI_CH3I_progress_completion_cond);
- /* Else, the condition shoul be satisfied to break the busy loop*/
}
# endif
diff --git a/src/mpid/common/thread/mpid_thread.h b/src/mpid/common/thread/mpid_thread.h
index df9c355..8386cf2 100644
--- a/src/mpid/common/thread/mpid_thread.h
+++ b/src/mpid/common/thread/mpid_thread.h
@@ -314,14 +314,6 @@ do { \
("mutex_unlock failed, err_=%d (%s)",err_,MPIU_Strerror(err_))); \
} while (0)
-#define MPID_Thread_mutex_lock_low(mutex_) \
-do { \
- int err_; \
- MPIU_Thread_mutex_lock_low((mutex_), &err_); \
- MPIU_Assert_fmt_msg(err_ == MPIU_THREAD_SUCCESS, \
- ("mutex_lock failed, err_=%d (%s)",err_,MPIU_Strerror(err_))); \
-} while (0)
-
#define MPID_Thread_mutex_trylock(mutex_, flag_) \
do { \
int err_; \
-----------------------------------------------------------------------
Summary of changes:
src/include/mpiimplthreadpost.h | 8 +-
src/include/thread/mpiu_thread_posix_funcs.h | 535 +++++-----------------
src/include/thread/mpiu_thread_posix_types.h | 53 +---
src/mpi/init/async.c | 32 +-
src/mpi/init/initthread.c | 35 --
src/mpid/ch3/channels/nemesis/src/ch3_progress.c | 15 +-
src/mpid/common/thread/mpid_thread.h | 8 -
7 files changed, 133 insertions(+), 553 deletions(-)
hooks/post-receive
--
MPICH primary repository
1
0
[mpich] MPICH primary repository branch, master, updated. v3.2b1-78-g17e31e5
by noreply@mpich.org 17 Apr '15
by noreply@mpich.org 17 Apr '15
17 Apr '15
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "MPICH primary repository".
The branch, master has been updated
via 17e31e59b491fc1383653ddd41f8264e8c5c69e6 (commit)
from 73e3211228c525b08221fbfcb161f88878f8cd24 (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
http://git.mpich.org/mpich.git/commitdiff/17e31e59b491fc1383653ddd41f8264e8…
commit 17e31e59b491fc1383653ddd41f8264e8c5c69e6
Author: Halim Amer <aamer(a)anl.gov>
Date: Fri Apr 17 19:17:23 2015 -0500
Applied the PPoPP patch
diff --git a/src/include/mpiimplthreadpost.h b/src/include/mpiimplthreadpost.h
index be4578f..60a0882 100644
--- a/src/include/mpiimplthreadpost.h
+++ b/src/include/mpiimplthreadpost.h
@@ -191,7 +191,13 @@ MPIU_Thread_CS_yield_lockname_recursive_impl_(enum MPIU_Nest_mutexes kind,
MPID_Thread_mutex_unlock(mutex);
MPID_Thread_yield();
- MPID_Thread_mutex_lock(mutex);
+
+/* Use low priority here because this thread has a lower probability
+ * to do useful work compared to others outside the progress engine.
+ * This is only effective with the priority lock.
+ * The other lock types implement a plain lock underneath.
+ */
+ MPID_Thread_mutex_lock_low(mutex);
}
/* undef for safety, this is a commonly-included header */
diff --git a/src/include/thread/mpiu_thread_posix_funcs.h b/src/include/thread/mpiu_thread_posix_funcs.h
index 10af9fc..facea60 100644
--- a/src/include/thread/mpiu_thread_posix_funcs.h
+++ b/src/include/thread/mpiu_thread_posix_funcs.h
@@ -8,8 +8,10 @@
/*
* Threads
*/
-
+#include <limits.h>
#include "mpiu_process_wrappers.h" /* for MPIU_PW_Sched_yield */
+#include "mpidbg.h"
+#include "opa_primitives.h"
/*
One of PTHREAD_MUTEX_RECURSIVE_NP and PTHREAD_MUTEX_RECURSIVE seem to be
@@ -53,144 +55,443 @@ do { \
MPIU_DBG_MSG(THREAD,VERBOSE,"exit MPIU_Thread_yield"); \
} while (0)
+/*----------------------*/
+/* Ticket Lock Routines */
+/*----------------------*/
-/*
- * Mutexes
- */
+static inline int ticket_lock_init(ticket_lock_t *lock)
+{
+ OPA_store_int(&lock->next_ticket, 0);
+ OPA_store_int(&lock->now_serving, 0);
+ return 0;
+}
+
+/* Atomically increment the nex_ticket counter and get my ticket.
+ Then spin on now_serving until it equals my ticket.
+ */
+static inline int ticket_acquire_lock(ticket_lock_t* lock)
+{
+ int my_ticket = OPA_fetch_and_add_int(&lock->next_ticket, 1);
+ while(OPA_load_int(&lock->now_serving) != my_ticket)
+ ;
+ return 0;
+}
+
+/* Release the lock
+ */
+static inline int ticket_release_lock(ticket_lock_t* lock)
+{
+ /* Avoid compiler reordering before releasing the lock*/
+ OPA_compiler_barrier();
+ OPA_incr_int(&lock->now_serving);
+ return 0;
+}
+
+/*------------------------*/
+/* Priority Lock Routines */
+/*------------------------*/
+
+static inline int priority_lock_init(priority_lock_t *lock)
+{
+ OPA_store_int(&lock->next_ticket_H, 0);
+ OPA_store_int(&lock->now_serving_H, 0);
+ OPA_store_int(&lock->next_ticket_L, 0);
+ OPA_store_int(&lock->now_serving_L, 0);
+ OPA_store_int(&lock->next_ticket_B, 0);
+ OPA_store_int(&lock->now_serving_B, 0);
+ lock->already_blocked = 0;
+ return 0;
+}
+
+/* First wait my turn in this priority level, Then if I am the first
+ one to block the LPRs, wait for the last LPR to terminate.
+ */
+static inline int priority_acquire_lock(priority_lock_t* lock)
+{
+ int my_ticket = OPA_fetch_and_add_int(&lock->next_ticket_H, 1);
+ while(OPA_load_int(&lock->now_serving_H) != my_ticket)
+ ;
+ int B_ticket;
+ if(!lock->already_blocked)
+ {
+ B_ticket = OPA_fetch_and_add_int(&lock->next_ticket_B, 1);
+ while(OPA_load_int(&lock->now_serving_B) != B_ticket)
+ ;
+ lock->already_blocked = 1;
+ }
+ lock->last_acquisition_priority = HIGH_PRIORITY;
+ return 0;
+}
+static inline int priority_acquire_lock_low(priority_lock_t* lock)
+{
+ int my_ticket = OPA_fetch_and_add_int(&lock->next_ticket_L, 1);
+ while(OPA_load_int(&lock->now_serving_L) != my_ticket)
+ ;
+ int B_ticket = OPA_fetch_and_add_int(&lock->next_ticket_B, 1);
+ while(OPA_load_int(&lock->now_serving_B) != B_ticket)
+ ;
+ lock->last_acquisition_priority = LOW_PRIORITY;
+ return 0;
+}
-/* FIXME: mutex creation and destruction should be implemented as routines
- because there is no reason to use macros (these are not on the performance
- critical path). Making these macros requires that any code that might use
- these must load all of the pthread.h (or other thread library) support.
+/* Release the lock
+ */
+static inline int priority_release_lock(priority_lock_t* lock)
+{
+ int err=0;
+ if(lock->last_acquisition_priority==HIGH_PRIORITY)
+ err = priority_release_lock_high(lock);
+ else
+ err = priority_release_lock_low(lock);
+ return err;
+}
+
+static inline int priority_release_lock_high(priority_lock_t* lock)
+{
+ /* Avoid compiler reordering before releasing the lock*/
+ OPA_compiler_barrier();
+ /* Only me in the HPRs queue -> let the LPRs pass */
+ if(OPA_load_int(&lock->now_serving_H) == OPA_load_int(&lock->next_ticket_H) - 1)
+ {
+ lock->already_blocked = 0;
+ OPA_incr_int(&lock->now_serving_B);
+ }
+ OPA_incr_int(&lock->now_serving_H);
+ return 0;
+}
+
+static inline int priority_release_lock_low(priority_lock_t* lock)
+{
+ /* Avoid compiler reordering before releasing the lock*/
+ OPA_compiler_barrier();
+ OPA_incr_int(&lock->now_serving_B);
+ OPA_incr_int(&lock->now_serving_L);
+ return 0;
+}
+
+/*
+ * MPIU Mutexes: encapsulate lower level locks like pthread mutexes
+ * and ticket locks
*/
-/* FIXME: using constant initializer if available */
#if !defined(MPICH_DEBUG_MUTEX) || !defined(PTHREAD_MUTEX_ERRORCHECK_VALUE)
-#define MPIU_Thread_mutex_create(mutex_ptr_, err_ptr_) \
-do { \
- int err__; \
- \
- err__ = pthread_mutex_init((mutex_ptr_), NULL); \
- /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
- *(int *)(err_ptr_) = err__; \
- MPIU_DBG_MSG_P(THREAD,TYPICAL,"Created MPIU_Thread_mutex %p", (mutex_ptr_)); \
-} while (0)
+static inline void MPIU_Thread_mutex_create(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
+{
+ int err__=0;
+ switch(MPIU_lock_type){
+ case MPIU_MUTEX:
+ {
+ err__ = pthread_mutex_init(&mutex_ptr_->pthread_lock, NULL);
+ break;
+ }
+ case MPIU_TICKET:
+ {
+ err__ = ticket_lock_init(&mutex_ptr_->ticket_lock);
+ break;
+ }
+ case MPIU_PRIORITY:
+ {
+ err__ = priority_lock_init(&mutex_ptr_->priority_lock);
+ break;
+ }
+ }
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */
+ *(int *)(err_ptr_) = err__;
+ MPIU_DBG_MSG_P(THREAD,TYPICAL,"Created MPIU_Thread_mutex %p", (mutex_ptr_));
+}
#else /* MPICH_DEBUG_MUTEX */
-#define MPIU_Thread_mutex_create(mutex_ptr_, err_ptr_) \
-do { \
- int err__; \
- pthread_mutexattr_t attr__; \
- \
- /* FIXME this used to be PTHREAD_MUTEX_ERRORCHECK_NP, but we had to change
- it for the thread granularity work when we needed recursive mutexes. We
- should go through this code and see if there's any good way to implement
- error checked versions with the recursive mutexes. */ \
- pthread_mutexattr_init(&attr__); \
- pthread_mutexattr_settype(&attr__, PTHREAD_MUTEX_ERRORCHECK_VALUE); \
- err__ = pthread_mutex_init((mutex_ptr_), &attr__); \
- if (err__) \
+static inline void MPIU_Thread_mutex_create(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
+{
+ int err__=0;
+ switch(MPIU_lock_type){
+ case MPIU_MUTEX:
+ {
+ pthread_mutexattr_t attr__;
+ /* FIXME this used to be PTHREAD_MUTEX_ERRORCHECK_NP, but we had to change
+ it for the thread granularity work when we needed recursive mutexes. We
+ should go through this code and see if there's any good way to implement
+ error checked versions with the recursive mutexes. */
+ pthread_mutexattr_init(&attr__);
+ pthread_mutexattr_settype(&attr__, PTHREAD_MUTEX_ERRORCHECK_VALUE);
+ err__ = pthread_mutex_init(&mutex_ptr_->pthread_lock, &attr__);
+ break;
+ }
+ case MPIU_TICKET:
+ {
+ err__ = ticket_lock_init(&mutex_ptr_->ticket_lock);
+ break;
+ }
+ case MPIU_PRIORITY:
+ {
+ err__ = priority_lock_init(&mutex_ptr_->priority_lock);
+ break;
+ }
+ }
+ if (err__)
MPIU_Internal_sys_error_printf("pthread_mutex_init", err__, \
- " %s:%d\n", __FILE__, __LINE__);\
- /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
- *(int *)(err_ptr_) = err__; \
- MPIU_DBG_MSG_P(THREAD,TYPICAL,"Created MPIU_Thread_mutex %p", (mutex_ptr_)); \
-} while (0)
+ " %s:%d\n", __FILE__, __LINE__);
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */
+ *(int *)(err_ptr_) = err__;
+ MPIU_DBG_MSG_P(THREAD,TYPICAL,"Created MPIU_Thread_mutex %p", (mutex_ptr_));
+}
#endif
-#define MPIU_Thread_mutex_destroy(mutex_ptr_, err_ptr_) \
-do { \
- int err__; \
- \
- MPIU_DBG_MSG_P(THREAD,TYPICAL,"About to destroy MPIU_Thread_mutex %p", (mutex_ptr_)); \
- err__ = pthread_mutex_destroy(mutex_ptr_); \
- /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
- *(int *)(err_ptr_) = err__; \
-} while (0)
+static inline void MPIU_Thread_mutex_destroy(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
+{
+ int err__=0;
+
+ MPIU_DBG_MSG_P(THREAD,TYPICAL,"About to destroy MPIU_Thread_mutex %p", (mutex_ptr_));
+
+ switch(MPIU_lock_type){
+ case MPIU_MUTEX:
+ {
+ err__ = pthread_mutex_destroy(&mutex_ptr_->pthread_lock);
+ break;
+ }
+ case MPIU_TICKET:
+ {
+ /* FIXME Should we do something here ?*/
+ err__ = 0;
+ break;
+ }
+ case MPIU_PRIORITY:
+ {
+ /* FIXME Should we do something here ?*/
+ err__ = 0;
+ break;
+ }
+ }
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */
+ *(int *)(err_ptr_) = err__;
+}
#ifndef MPICH_DEBUG_MUTEX
-#define MPIU_Thread_mutex_lock(mutex_ptr_, err_ptr_) \
-do { \
- int err__; \
- MPIU_DBG_MSG_P(THREAD,VERBOSE,"enter MPIU_Thread_mutex_lock %p", (mutex_ptr_)); \
- err__ = pthread_mutex_lock(mutex_ptr_); \
- /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
- *(int *)(err_ptr_) = err__; \
- MPIU_DBG_MSG_P(THREAD,VERBOSE,"exit MPIU_Thread_mutex_lock %p", (mutex_ptr_)); \
-} while (0)
+static inline void MPIU_Thread_mutex_lock(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
+{
+ int err__=0;
+ MPIU_DBG_MSG_P(THREAD,VERBOSE,"enter MPIU_Thread_mutex_lock %p", (mutex_ptr_));
+
+ switch(MPIU_lock_type){
+ case MPIU_MUTEX:
+ {
+ err__ = pthread_mutex_lock(&mutex_ptr_->pthread_lock);
+ break;
+ }
+ case MPIU_TICKET:
+ {
+ err__ = ticket_acquire_lock(&mutex_ptr_->ticket_lock);
+ break;
+ }
+ case MPIU_PRIORITY:
+ {
+ err__ = priority_acquire_lock(&mutex_ptr_->priority_lock);
+ break;
+ }
+ }
+
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */
+ *(int *)(err_ptr_) = err__;
+ MPIU_DBG_MSG_P(THREAD,VERBOSE,"exit MPIU_Thread_mutex_lock %p", (mutex_ptr_));
+}
#else /* MPICH_DEBUG_MUTEX */
-#define MPIU_Thread_mutex_lock(mutex_ptr_, err_ptr_) \
-do { \
- int err__; \
- MPIU_DBG_MSG_P(THREAD,VERBOSE,"enter MPIU_Thread_mutex_lock %p", (mutex_ptr_)); \
- err__ = pthread_mutex_lock(mutex_ptr_); \
- if (err__) \
- { \
- MPIU_DBG_MSG_S(THREAD,TERSE," mutex lock error: %s", MPIU_Strerror(err__)); \
+static inline void MPIU_Thread_mutex_lock(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
+{
+ int err__=0;
+ MPIU_DBG_MSG_P(THREAD,VERBOSE,"enter MPIU_Thread_mutex_lock %p", (mutex_ptr_));
+ switch(MPIU_lock_type){
+ case MPIU_MUTEX:
+ {
+ err__ = pthread_mutex_lock(&mutex_ptr_->pthread_lock);
+ break;
+ }
+ case MPIU_TICKET:
+ {
+ err__ = ticket_acquire_lock(&mutex_ptr_->ticket_lock);
+ break;
+ }
+ case MPIU_PRIORITY:
+ {
+ err__ = priority_acquire_lock(&mutex_ptr_->priority_lock);
+ break;
+ }
+ }
+
+ if (err__)
+ {
+ MPIU_DBG_MSG_S(THREAD,TERSE," mutex lock error: %s", MPIU_Strerror(err__));
MPIU_Internal_sys_error_printf("pthread_mutex_lock", err__,\
- " %s:%d\n", __FILE__, __LINE__);\
- } \
- /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
- *(int *)(err_ptr_) = err__; \
- MPIU_DBG_MSG_P(THREAD,VERBOSE,"exit MPIU_Thread_mutex_lock %p", (mutex_ptr_)); \
-} while (0)
+ " %s:%d\n", __FILE__, __LINE__);
+ }
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */
+ *(int *)(err_ptr_) = err__;
+ MPIU_DBG_MSG_P(THREAD,VERBOSE,"exit MPIU_Thread_mutex_lock %p", (mutex_ptr_));
+}
#endif
#ifndef MPICH_DEBUG_MUTEX
-#define MPIU_Thread_mutex_unlock(mutex_ptr_, err_ptr_) \
-do { \
- int err__; \
- \
- MPIU_DBG_MSG_P(THREAD,TYPICAL,"MPIU_Thread_mutex_unlock %p", (mutex_ptr_)); \
- err__ = pthread_mutex_unlock(mutex_ptr_); \
- /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
- *(int *)(err_ptr_) = err__; \
-} while (0)
+static inline void MPIU_Thread_mutex_lock_low(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
+{
+ int err__=0;
+ MPIU_DBG_MSG_P(THREAD,VERBOSE,"enter MPIU_Thread_mutex_lock %p", (mutex_ptr_));
+
+ switch(MPIU_lock_type){
+ case MPIU_MUTEX:
+ {
+ err__ = pthread_mutex_lock(&mutex_ptr_->pthread_lock);
+ break;
+ }
+ case MPIU_TICKET:
+ {
+ err__ = ticket_acquire_lock(&mutex_ptr_->ticket_lock);
+ break;
+ }
+ case MPIU_PRIORITY:
+ {
+ err__ = priority_acquire_lock_low(&mutex_ptr_->priority_lock);
+ break;
+ }
+ }
+
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */
+ *(int *)(err_ptr_) = err__;
+ MPIU_DBG_MSG_P(THREAD,VERBOSE,"exit MPIU_Thread_mutex_lock %p", (mutex_ptr_));
+}
#else /* MPICH_DEBUG_MUTEX */
-#define MPIU_Thread_mutex_unlock(mutex_ptr_, err_ptr_) \
-do { \
- int err__; \
- \
- MPIU_DBG_MSG_P(THREAD,VERBOSE,"MPIU_Thread_mutex_unlock %p", (mutex_ptr_)); \
- err__ = pthread_mutex_unlock(mutex_ptr_); \
- if (err__) \
- { \
- MPIU_DBG_MSG_S(THREAD,TERSE," mutex unlock error: %s", MPIU_Strerror(err__)); \
- MPIU_Internal_sys_error_printf("pthread_mutex_unlock", err__, \
- " %s:%d\n", __FILE__, __LINE__); \
- } \
- /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
- *(int *)(err_ptr_) = err__; \
-} while (0)
+static inline void MPIU_Thread_mutex_lock_low(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
+{
+ int err__=0;
+ MPIU_DBG_MSG_P(THREAD,VERBOSE,"enter MPIU_Thread_mutex_lock %p", (mutex_ptr_));
+ switch(MPIU_lock_type){
+ case MPIU_MUTEX:
+ {
+ err__ = pthread_mutex_lock(&mutex_ptr_->pthread_lock);
+ break;
+ }
+ case MPIU_TICKET:
+ {
+ err__ = ticket_acquire_lock(&mutex_ptr_->ticket_lock);
+ break;
+ }
+ case MPIU_PRIORITY:
+ {
+ err__ = priority_acquire_lock_low(&mutex_ptr_->priority_lock);
+ break;
+ }
+ }
+
+ if (err__)
+ {
+ MPIU_DBG_MSG_S(THREAD,TERSE," mutex lock error: %s", MPIU_Strerror(err__));
+ MPIU_Internal_sys_error_printf("pthread_mutex_lock", err__,\
+ " %s:%d\n", __FILE__, __LINE__);
+ }
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */
+ *(int *)(err_ptr_) = err__;
+ MPIU_DBG_MSG_P(THREAD,VERBOSE,"exit MPIU_Thread_mutex_lock %p", (mutex_ptr_));
+}
#endif
#ifndef MPICH_DEBUG_MUTEX
-#define MPIU_Thread_mutex_trylock(mutex_ptr_, flag_ptr_, err_ptr_) \
-do { \
- int err__; \
- \
- err__ = pthread_mutex_trylock(mutex_ptr_); \
- *(flag_ptr_) = (err__ == 0) ? TRUE : FALSE; \
- MPIU_DBG_MSG_FMT(THREAD,VERBOSE,(MPIU_DBG_FDEST, "MPIU_Thread_mutex_trylock mutex=%p result=%s", (mutex_ptr_), (*(flag_ptr_) ? "success" : "failure"))); \
- *(int *)(err_ptr_) = (err__ == EBUSY) ? MPIU_THREAD_SUCCESS : err__; \
- /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
-} while (0)
+static inline void MPIU_Thread_mutex_unlock(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
+{
+ int err__=0;
+
+ MPIU_DBG_MSG_P(THREAD,TYPICAL,"MPIU_Thread_mutex_unlock %p", (mutex_ptr_));
+
+ switch(MPIU_lock_type){
+ case MPIU_MUTEX:
+ {
+ err__ = pthread_mutex_unlock(&mutex_ptr_->pthread_lock);
+ break;
+ }
+ case MPIU_TICKET:
+ {
+ err__ = ticket_release_lock(&mutex_ptr_->ticket_lock);
+ break;
+ }
+ case MPIU_PRIORITY:
+ {
+ err__ = priority_release_lock(&mutex_ptr_->priority_lock);
+ break;
+ }
+ }
+
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */
+ *(int *)(err_ptr_) = err__;
+}
#else /* MPICH_DEBUG_MUTEX */
-#define MPIU_Thread_mutex_trylock(mutex_ptr_, flag_ptr_, err_ptr_) \
-do { \
- int err__; \
- \
- err__ = pthread_mutex_trylock(mutex_ptr_); \
- if (err__ && err__ != EBUSY) \
- { \
- MPIU_DBG_MSG_S(THREAD,TERSE," mutex trylock error: %s", MPIU_Strerror(err__)); \
+static inline void MPIU_Thread_mutex_unlock(MPIU_Thread_mutex_t* mutex_ptr_, int* err_ptr_)
+{
+ int err__=0;
+
+ MPIU_DBG_MSG_P(THREAD,VERBOSE,"MPIU_Thread_mutex_unlock %p", (mutex_ptr_));
+
+ switch(MPIU_lock_type){
+ case MPIU_MUTEX:
+ {
+ err__ = pthread_mutex_unlock(&mutex_ptr_->pthread_lock);
+ break;
+ }
+ case MPIU_TICKET:
+ {
+ err__ = ticket_release_lock(&mutex_ptr_->ticket_lock);
+ break;
+ }
+ case MPIU_PRIORITY:
+ {
+ err__ = priority_release_lock(&mutex_ptr_->priority_lock);
+ break;
+ }
+ }
+ if (err__)
+ {
+ MPIU_DBG_MSG_S(THREAD,TERSE," mutex unlock error: %s", MPIU_Strerror(err__));
+ MPIU_Internal_sys_error_printf("pthread_mutex_unlock", err__,\
+ " %s:%d\n", __FILE__, __LINE__);
+ }
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */
+ *(int *)(err_ptr_) = err__;
+}
+#endif
+
+#ifndef MPICH_DEBUG_MUTEX
+static inline void MPIU_Thread_mutex_trylock(MPIU_Thread_mutex_t* mutex_ptr_, int* flag_ptr_, int* err_ptr_)
+{
+ int err__=0;
+
+ if(MPIU_lock_type == MPIU_MUTEX)
+ err__ = pthread_mutex_trylock(&mutex_ptr_->pthread_lock);
+ else
+ /*No trylock routine for ticket-based locks*/
+ err__ = 1;
+
+ *(flag_ptr_) = (err__ == 0) ? TRUE : FALSE;
+ MPIU_DBG_MSG_FMT(THREAD,VERBOSE,(MPIU_DBG_FDEST, "MPIU_Thread_mutex_trylock mutex=%p result=%s", (mutex_ptr_), (*(flag_ptr_) ? "success" : "failure")));
+ *(int *)(err_ptr_) = (err__ == EBUSY) ? MPIU_THREAD_SUCCESS : err__;
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */
+}
+#else /* MPICH_DEBUG_MUTEX */
+static inline void MPIU_Thread_mutex_trylock(MPIU_Thread_mutex_t* mutex_ptr_, int* flag_ptr_, int* err_ptr_)
+{
+ int err__=0;
+
+ if(MPIU_lock_type == MPIU_MUTEX)
+ err__ = pthread_mutex_trylock(&mutex_ptr_->pthread_lock);
+ else
+ /*No trylock routine for ticket-based locks*/
+ err__ = 1;
+
+ if (err__ && err__ != EBUSY)
+ {
+ MPIU_DBG_MSG_S(THREAD,TERSE," mutex trylock error: %s", MPIU_Strerror(err__));
MPIU_Internal_sys_error_printf("pthread_mutex_trylock", err__,\
- " %s:%d\n", __FILE__, __LINE__);\
- } \
- *(flag_ptr_) = (err__ == 0) ? TRUE : FALSE; \
- MPIU_DBG_MSG_FMT(THREAD,VERBOSE,(MPIU_DBG_FDEST, "MPIU_Thread_mutex_trylock mutex=%p result=%s", (mutex_ptr_), (*(flag_ptr_) ? "success" : "failure"))); \
- *(int *)(err_ptr_) = (err__ == EBUSY) ? MPIU_THREAD_SUCCESS : err__; \
- /* FIXME: convert error to an MPIU_THREAD_ERR value */ \
-} while (0)
+ " %s:%d\n", __FILE__, __LINE__);
+ }
+ *(flag_ptr_) = (err__ == 0) ? TRUE : FALSE;
+ MPIU_DBG_MSG_FMT(THREAD,VERBOSE,(MPIU_DBG_FDEST, "MPIU_Thread_mutex_trylock mutex=%p result=%s", (mutex_ptr_), (*(flag_ptr_) ? "success" : "failure")));
+ *(int *)(err_ptr_) = (err__ == EBUSY) ? MPIU_THREAD_SUCCESS : err__;
+ /* FIXME: convert error to an MPIU_THREAD_ERR value */
+}
#endif
/*
diff --git a/src/include/thread/mpiu_thread_posix_types.h b/src/include/thread/mpiu_thread_posix_types.h
index b38f7a1..c163b67 100644
--- a/src/include/thread/mpiu_thread_posix_types.h
+++ b/src/include/thread/mpiu_thread_posix_types.h
@@ -7,8 +7,59 @@
#include <errno.h>
#include <pthread.h>
+#include "opa_primitives.h"
+
+/* Define lock types here */
+typedef enum{
+ MPIU_MUTEX,
+ MPIU_TICKET,
+ MPIU_PRIORITY
+}MPIU_Thread_lock_impl_t;
+
+MPIU_Thread_lock_impl_t MPIU_lock_type;
+
+/*----------------------------*/
+/* Ticket lock data structure */
+/*----------------------------*/
+
+/* Define the lock as a structure of two counters */
+typedef struct ticket_lock_t{
+ OPA_int_t next_ticket;
+ OPA_int_t now_serving;
+} ticket_lock_t;
+
+/*------------------------------*/
+/* Priority lock data structure */
+/*------------------------------*/
+
+#define HIGH_PRIORITY 1
+#define LOW_PRIORITY 2
+
+ /* We only define for now 2 levels of priority: */
+ /* The high priority requests (HPRs) */
+ /* The low priority requests (LPRs) */
+typedef struct priority_lock_t{
+ OPA_int_t next_ticket_H __attribute__((aligned(64)));
+ OPA_int_t now_serving_H;
+ OPA_int_t next_ticket_L __attribute__((aligned(64)));
+ OPA_int_t now_serving_L;
+ /* In addition we include two other counters */
+ /* so high priority requests block the lower ones */
+ OPA_int_t next_ticket_B __attribute__((aligned(64)));
+ OPA_int_t now_serving_B;
+ /* This is to allow high priority requests know */
+ /* that low priority requests are already blocked */
+ /* by another high*/
+ unsigned already_blocked;
+ unsigned last_acquisition_priority;
+} priority_lock_t;
+
+typedef struct MPIU_Thread_mutex_t{
+ pthread_mutex_t pthread_lock;
+ ticket_lock_t ticket_lock;
+ priority_lock_t priority_lock;
+} MPIU_Thread_mutex_t;
-typedef pthread_mutex_t MPIU_Thread_mutex_t;
typedef pthread_cond_t MPIU_Thread_cond_t;
typedef pthread_t MPIU_Thread_id_t;
typedef pthread_key_t MPIU_Thread_tls_t;
diff --git a/src/mpi/init/async.c b/src/mpi/init/async.c
index c003827..a95e648 100644
--- a/src/mpi/init/async.c
+++ b/src/mpi/init/async.c
@@ -61,9 +61,10 @@ static void progress_fn(void * data)
MPIU_Thread_mutex_unlock(&progress_mutex, &mpi_errno);
MPIU_Assert(!mpi_errno);
-
- MPIU_Thread_cond_signal(&progress_cond, &mpi_errno);
- MPIU_Assert(!mpi_errno);
+ if(MPIU_lock_type==MPIU_MUTEX){
+ MPIU_Thread_cond_signal(&progress_cond, &mpi_errno);
+ MPIU_Assert(!mpi_errno);
+ } /* Else: busy loop will automatically break*/
MPIU_THREAD_CS_EXIT(ALLFUNC,);
@@ -137,16 +138,21 @@ int MPIR_Finalize_async_thread(void)
/* XXX DJG why is this unlock/lock necessary? Should we just YIELD here or later? */
MPIU_THREAD_CS_EXIT(ALLFUNC,);
- MPIU_Thread_mutex_lock(&progress_mutex, &mpi_errno);
- MPIU_Assert(!mpi_errno);
-
- while (!progress_thread_done) {
- MPIU_Thread_cond_wait(&progress_cond, &progress_mutex, &mpi_errno);
- MPIU_Assert(!mpi_errno);
- }
-
- MPIU_Thread_mutex_unlock(&progress_mutex, &mpi_errno);
- MPIU_Assert(!mpi_errno);
+ if(MPIU_lock_type==MPIU_MUTEX){
+ MPIU_Thread_mutex_lock(&progress_mutex, &mpi_errno);
+ MPIU_Assert(!mpi_errno);
+
+ while (!progress_thread_done) {
+ MPIU_Thread_cond_wait(&progress_cond, &progress_mutex.pthread_lock, &mpi_errno);
+ MPIU_Assert(!mpi_errno);
+ }
+
+ MPIU_Thread_mutex_unlock(&progress_mutex, &mpi_errno);
+ MPIU_Assert(!mpi_errno);
+ }
+ else
+ while (!progress_thread_done) ; /* busy loop */
+ /* No need to unlock the mutex */
mpi_errno = MPIR_Comm_free_impl(progress_comm_ptr);
MPIU_Assert(!mpi_errno);
diff --git a/src/mpi/init/initthread.c b/src/mpi/init/initthread.c
index aaf20c8..2b5db44 100644
--- a/src/mpi/init/initthread.c
+++ b/src/mpi/init/initthread.c
@@ -190,6 +190,16 @@ static int MPIR_Thread_CS_Init( void )
int err;
MPIU_THREADPRIV_DECL;
+ MPIU_lock_type = MPIU_MUTEX;
+ char *s;
+ s = getenv( "MPICH_LOCK_TYPE" );
+ if(s){
+ if(strcmp( "ticket", s ) == 0)
+ MPIU_lock_type = MPIU_TICKET;
+ else if (strcmp( "priority", s ) == 0)
+ MPIU_lock_type = MPIU_PRIORITY;
+ }
+
MPIU_Assert(MPICH_MAX_LOCKS >= MPIU_Nest_NUM_MUTEXES);
/* we create this at all granularities right now */
@@ -558,6 +568,30 @@ int MPIR_Init_thread(int * argc, char ***argv, int required, int * provided)
if (mpi_errno == MPI_SUCCESS)
mpi_errno = MPID_InitCompleted();
+ const char *s;
+
+ switch(MPIU_lock_type){
+ case MPIU_MUTEX:
+ {
+ s = "mutex";
+ break;
+ }
+ case MPIU_TICKET:
+ {
+ s = "ticket";
+ break;
+ }
+ case MPIU_PRIORITY:
+ {
+ s = "priority";
+ break;
+ }
+ }
+
+ if(MPIR_Process.comm_world->rank==0){
+ printf("\n[MPICH INFO] Critical section(s) based on %s \n\n", s);
+ }
+
fn_exit:
MPIU_THREAD_CS_EXIT(INIT,required);
/* Make fields of MPIR_Process global visible and set mpich_state
@@ -626,6 +660,7 @@ int MPI_Init_thread( int *argc, char ***argv, int required, int *provided )
{
int mpi_errno = MPI_SUCCESS;
int rc ATTRIBUTE((unused)), reqd = required;
+
MPID_MPI_INIT_STATE_DECL(MPID_STATE_MPI_INIT_THREAD);
rc = MPID_Wtime_init();
diff --git a/src/mpid/ch3/channels/nemesis/src/ch3_progress.c b/src/mpid/ch3/channels/nemesis/src/ch3_progress.c
index f022651..8dda2cd 100644
--- a/src/mpid/ch3/channels/nemesis/src/ch3_progress.c
+++ b/src/mpid/ch3/channels/nemesis/src/ch3_progress.c
@@ -563,13 +563,24 @@ static int MPIDI_CH3I_Progress_delay(unsigned int completion_count)
/* FIXME should be appropriately abstracted somehow */
# if defined(MPICH_IS_THREADED) && (MPIU_THREAD_GRANULARITY == MPIU_THREAD_GRANULARITY_GLOBAL)
{
+ if(MPIU_lock_type==MPIU_MUTEX){
while (1)
{
if (completion_count != OPA_load_int(&MPIDI_CH3I_progress_completion_count) ||
MPIDI_CH3I_progress_blocked != TRUE)
break;
- MPID_Thread_cond_wait(&MPIDI_CH3I_progress_completion_cond, &MPIR_ThreadInfo.global_mutex/*MPIDCOMM*/);
+ MPID_Thread_cond_wait(&MPIDI_CH3I_progress_completion_cond, &MPIR_ThreadInfo.global_mutex.pthread_lock/*MPIDCOMM*/);
}
+ }
+ else{
+ /* First release the lock and then enter a busy loop*/
+ MPID_Thread_mutex_unlock(&MPIR_ThreadInfo.global_mutex/*MPIDCOMM*/);
+ while (completion_count == OPA_load_int(&MPIDI_CH3I_progress_completion_count) &&
+ MPIDI_CH3I_progress_blocked == TRUE)
+ ; /* wait in a busy loop*/
+ /* Hold the lock again*/
+ MPID_Thread_mutex_lock(&MPIR_ThreadInfo.global_mutex/*MPIDCOMM*/);
+ }
}
# endif
@@ -593,7 +604,9 @@ static int MPIDI_CH3I_Progress_continue(unsigned int completion_count/*unused*/)
# if defined(MPICH_IS_THREADED) && (MPIU_THREAD_GRANULARITY == MPIU_THREAD_GRANULARITY_GLOBAL)
{
/* we currently hold the MPIDCOMM CS */
+ if(MPIU_lock_type==MPIU_MUTEX)
MPID_Thread_cond_broadcast(&MPIDI_CH3I_progress_completion_cond);
+ /* Else, the condition shoul be satisfied to break the busy loop*/
}
# endif
diff --git a/src/mpid/common/thread/mpid_thread.h b/src/mpid/common/thread/mpid_thread.h
index 8386cf2..df9c355 100644
--- a/src/mpid/common/thread/mpid_thread.h
+++ b/src/mpid/common/thread/mpid_thread.h
@@ -314,6 +314,14 @@ do { \
("mutex_unlock failed, err_=%d (%s)",err_,MPIU_Strerror(err_))); \
} while (0)
+#define MPID_Thread_mutex_lock_low(mutex_) \
+do { \
+ int err_; \
+ MPIU_Thread_mutex_lock_low((mutex_), &err_); \
+ MPIU_Assert_fmt_msg(err_ == MPIU_THREAD_SUCCESS, \
+ ("mutex_lock failed, err_=%d (%s)",err_,MPIU_Strerror(err_))); \
+} while (0)
+
#define MPID_Thread_mutex_trylock(mutex_, flag_) \
do { \
int err_; \
-----------------------------------------------------------------------
Summary of changes:
src/include/mpiimplthreadpost.h | 8 +-
src/include/thread/mpiu_thread_posix_funcs.h | 535 +++++++++++++++++-----
src/include/thread/mpiu_thread_posix_types.h | 53 +++-
src/mpi/init/async.c | 32 +-
src/mpi/init/initthread.c | 35 ++
src/mpid/ch3/channels/nemesis/src/ch3_progress.c | 15 +-
src/mpid/common/thread/mpid_thread.h | 8 +
7 files changed, 553 insertions(+), 133 deletions(-)
hooks/post-receive
--
MPICH primary repository
1
0
[mpich] MPICH primary repository branch, master, updated. v3.2b1-77-g73e3211
by noreply@mpich.org 17 Apr '15
by noreply@mpich.org 17 Apr '15
17 Apr '15
This is an automated email from the git hooks/post-receive script. It was
generated because a ref change was pushed to the repository containing
the project "MPICH primary repository".
The branch, master has been updated
via 73e3211228c525b08221fbfcb161f88878f8cd24 (commit)
from 4309ba574293e7bdf2ec03224e4afa948a1e6fd9 (commit)
Those revisions listed above that are new to this repository have
not appeared on any other notification email; so we list those
revisions in full, below.
- Log -----------------------------------------------------------------
http://git.mpich.org/mpich.git/commitdiff/73e3211228c525b08221fbfcb161f8887…
commit 73e3211228c525b08221fbfcb161f88878f8cd24
Author: Ken Raffenetti <raffenet(a)mcs.anl.gov>
Date: Thu Apr 16 14:10:40 2015 -0500
portals4: fix anysource_matched
This fix, along with a pending patch to the Portal4 reference implementation,
should make anysource_matched a more reliable operation for multithreaded apps. We
were seeing a race condition where an ME would unlink successfully, but an
event matching it would still arrive in the queue. CH3 can now reliably search
the netmod queue for matched MPI_ANY_SOURCE requests.
The reason that we no longer assert that an MPI_ANY_SOURCE request was removed
from the CH3 queue is that FDP (find and dequeue posted) operations will remove
the request from the queue, if it is known to be already matched by the netmod,
even if it has not yet completed.
Fixes #2199
Signed-off-by: Antonio J. Pena <apenya(a)mcs.anl.gov>
diff --git a/src/mpid/ch3/channels/nemesis/netmod/portals4/ptl_recv.c b/src/mpid/ch3/channels/nemesis/netmod/portals4/ptl_recv.c
index ec6d90a..f732751 100644
--- a/src/mpid/ch3/channels/nemesis/netmod/portals4/ptl_recv.c
+++ b/src/mpid/ch3/channels/nemesis/netmod/portals4/ptl_recv.c
@@ -22,7 +22,10 @@ static void dequeue_req(const ptl_event_t *e)
REQ_PTL(rreq)->put_me = PTL_INVALID_HANDLE;
found = MPIDI_CH3U_Recvq_DP(rreq);
- MPIU_Assert(found);
+ /* an MPI_ANY_SOURCE request may have been previously removed from the
+ CH3 queue by an FDP (find and dequeue posted) operation */
+ if (rreq->dev.match.parts.rank != MPI_ANY_SOURCE)
+ MPIU_Assert(found);
rreq->status.MPI_ERROR = MPI_SUCCESS;
rreq->status.MPI_SOURCE = NPTL_MATCH_GET_RANK(e->match_bits);
@@ -597,15 +600,12 @@ static int cancel_recv(MPID_Request *rreq, int *cancelled)
/* An invalid handle indicates the operation has been completed
and the matching list entry unlinked. At that point, the operation
cannot be cancelled. */
- if (REQ_PTL(rreq)->put_me != PTL_INVALID_HANDLE) {
- ptl_err = PtlMEUnlink(REQ_PTL(rreq)->put_me);
- if (ptl_err == PTL_OK)
- *cancelled = TRUE;
- /* FIXME: if we properly invalidate matching list entry handles, we should be
- able to ensure an unlink operation results in either PTL_OK or PTL_IN_USE.
- Anything else would be an error. For now, though, we assume anything but PTL_OK
- is uncancelable and return. */
- }
+ if (REQ_PTL(rreq)->put_me == PTL_INVALID_HANDLE)
+ goto fn_exit;
+
+ ptl_err = PtlMEUnlink(REQ_PTL(rreq)->put_me);
+ if (ptl_err == PTL_OK)
+ *cancelled = TRUE;
fn_exit:
MPIDI_FUNC_EXIT(MPID_STATE_CANCEL_RECV);
@@ -635,7 +635,7 @@ int MPID_nem_ptl_anysource_matched(MPID_Request *rreq)
fn_exit:
MPIDI_FUNC_EXIT(MPID_STATE_MPID_NEM_PTL_ANYSOURCE_MATCHED);
- return MPI_SUCCESS;
+ return !cancelled;
fn_fail:
goto fn_exit;
}
diff --git a/test/mpi/threads/pt2pt/testlist b/test/mpi/threads/pt2pt/testlist
index 08e25a4..b6ab43e 100644
--- a/test/mpi/threads/pt2pt/testlist
+++ b/test/mpi/threads/pt2pt/testlist
@@ -1,6 +1,6 @@
threads 2 timeLimit=600
threaded_sr 2
-alltoall 4 xfail=ticket2199
+alltoall 4
sendselfth 1
multisend 2
multisend2 5
-----------------------------------------------------------------------
Summary of changes:
.../channels/nemesis/netmod/portals4/ptl_recv.c | 22 ++++++++++----------
test/mpi/threads/pt2pt/testlist | 2 +-
2 files changed, 12 insertions(+), 12 deletions(-)
hooks/post-receive
--
MPICH primary repository
1
0