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, get-acc-from-sameer-3 has been created at e49ee8c03d72671d41616a07f818023e69dae365 (commit) - Log ----------------------------------------------------------------- http://git.mpich.org/mpich.git/commitdiff/e49ee8c03d72671d41616a07f818023e69... commit e49ee8c03d72671d41616a07f818023e69dae365 Author: Sameer Kumar <[email protected]> Date: Wed Dec 18 04:40:40 2013 -0600 Implement MPID_Get_accumulate() Does not include the win request tls optimization Signed-off-by: Michael Blocksome <[email protected]> diff --git a/src/mpid/pamid/include/mpidi_datatypes.h b/src/mpid/pamid/include/mpidi_datatypes.h index 8dcd0b7..1b4dee8 100644 --- a/src/mpid/pamid/include/mpidi_datatypes.h +++ b/src/mpid/pamid/include/mpidi_datatypes.h @@ -158,6 +158,8 @@ enum MPIDI_Protocols_WinCtrl, MPIDI_Protocols_WinAccum, MPIDI_Protocols_RVZ_zerobyte, + MPIDI_Protocols_WinGetAccum, + MPIDI_Protocols_WinGetAccumAck, #ifdef DYNAMIC_TASKING MPIDI_Protocols_Dyntask, MPIDI_Protocols_Dyntask_disconnect, diff --git a/src/mpid/pamid/include/mpidi_prototypes.h b/src/mpid/pamid/include/mpidi_prototypes.h index ec50f71..dfdd1b8 100644 --- a/src/mpid/pamid/include/mpidi_prototypes.h +++ b/src/mpid/pamid/include/mpidi_prototypes.h @@ -204,6 +204,25 @@ MPIDI_WinControlCB(pami_context_t context, pami_endpoint_t sender, pami_recv_t * recv); +void +MPIDI_WinGetAccumCB(pami_context_t context, + void * cookie, + const void * _control, + size_t size, + const void * sndbuf, + size_t sndlen, + pami_endpoint_t sender, + pami_recv_t * recv); +void +MPIDI_WinGetAccumAckCB(pami_context_t context, + void * cookie, + const void * _control, + size_t size, + const void * sndbuf, + size_t sndlen, + pami_endpoint_t sender, + pami_recv_t * recv); + /** \brief Helper function to complete a rendevous transfer */ pami_result_t MPIDI_RendezvousTransfer(pami_context_t context, void* rreq); pami_result_t MPIDI_RendezvousTransfer_SyncAck(pami_context_t context, void* rreq); diff --git a/src/mpid/pamid/src/misc/mpid_unimpl.c b/src/mpid/pamid/src/misc/mpid_unimpl.c index cef7096..03e9783 100644 --- a/src/mpid/pamid/src/misc/mpid_unimpl.c +++ b/src/mpid/pamid/src/misc/mpid_unimpl.c @@ -161,12 +161,4 @@ int MPID_Rget(void *origin_addr, int origin_count, return 0; } -int MPID_Get_accumulate(const void *origin_addr, int origin_count, - MPI_Datatype origin_datatype, void *result_addr, int result_count, - MPI_Datatype result_datatype, int target_rank, MPI_Aint target_disp, - int target_count, MPI_Datatype target_datatype, MPI_Op op, MPID_Win *win) -{ - MPID_abort(); - return 0; -} diff --git a/src/mpid/pamid/src/mpid_init.c b/src/mpid/pamid/src/mpid_init.c index 445b33d..5809ec1 100644 --- a/src/mpid/pamid/src/mpid_init.c +++ b/src/mpid/pamid/src/mpid_init.c @@ -151,6 +151,8 @@ static struct struct protocol_t WinCtrl; struct protocol_t WinAccum; struct protocol_t RVZ_zerobyte; + struct protocol_t WinGetAccum; + struct protocol_t WinGetAccumAck; #ifdef DYNAMIC_TASKING struct protocol_t Dyntask; struct protocol_t Dyntask_disconnect; @@ -252,6 +254,26 @@ static struct }, .immediate_min = sizeof(MPIDI_MsgEnvelope), }, + .WinGetAccum = { + .func = MPIDI_WinGetAccumCB, + .dispatch = MPIDI_Protocols_WinGetAccum, + .options = { + .consistency = PAMI_HINT_ENABLE, + .long_header = PAMI_HINT_DISABLE, + .recv_immediate = PAMI_HINT_DISABLE, + }, + .immediate_min = sizeof(MPIDI_Win_GetAccMsgInfo), + }, + .WinGetAccumAck = { + .func = MPIDI_WinGetAccumAckCB, + .dispatch = MPIDI_Protocols_WinGetAccumAck, + .options = { + .consistency = PAMI_HINT_ENABLE, + .long_header = PAMI_HINT_DISABLE, + .recv_immediate = PAMI_HINT_DISABLE, + }, + .immediate_min = sizeof(MPIDI_Win_GetAccMsgInfo), + }, #ifdef DYNAMIC_TASKING .Dyntask = { .func = MPIDI_Recvfrom_remote_world, @@ -802,6 +824,8 @@ MPIDI_PAMI_dispath_init() MPIDI_PAMI_dispath_set(MPIDI_Protocols_WinCtrl, &proto_list.WinCtrl, NULL); MPIDI_PAMI_dispath_set(MPIDI_Protocols_WinAccum, &proto_list.WinAccum, NULL); MPIDI_PAMI_dispath_set(MPIDI_Protocols_RVZ_zerobyte, &proto_list.RVZ_zerobyte, NULL); + MPIDI_PAMI_dispath_set(MPIDI_Protocols_WinGetAccum, &proto_list.WinGetAccum, NULL); + MPIDI_PAMI_dispath_set(MPIDI_Protocols_WinGetAccumAck, &proto_list.WinGetAccumAck, NULL); #ifdef DYNAMIC_TASKING MPIDI_PAMI_dispath_set(MPIDI_Protocols_Dyntask, &proto_list.Dyntask, NULL); MPIDI_PAMI_dispath_set(MPIDI_Protocols_Dyntask_disconnect, &proto_list.Dyntask_disconnect, NULL); diff --git a/src/mpid/pamid/src/onesided/Makefile.mk b/src/mpid/pamid/src/onesided/Makefile.mk index 843f06d..901e6ab 100644 --- a/src/mpid/pamid/src/onesided/Makefile.mk +++ b/src/mpid/pamid/src/onesided/Makefile.mk @@ -28,6 +28,7 @@ noinst_HEADERS += \ lib_lib@MPILIBNAME@_la_SOURCES += \ src/mpid/pamid/src/onesided/mpid_1s.c \ src/mpid/pamid/src/onesided/mpid_win_accumulate.c \ + src/mpid/pamid/src/onesided/mpid_win_get_accumulate.c \ src/mpid/pamid/src/onesided/mpid_win_create.c \ src/mpid/pamid/src/onesided/mpid_win_fence.c \ src/mpid/pamid/src/onesided/mpid_win_free.c \ diff --git a/src/mpid/pamid/src/onesided/mpid_win_accumulate.c b/src/mpid/pamid/src/onesided/mpid_win_accumulate.c index c2878d8..ecdcf14 100644 --- a/src/mpid/pamid/src/onesided/mpid_win_accumulate.c +++ b/src/mpid/pamid/src/onesided/mpid_win_accumulate.c @@ -95,7 +95,7 @@ MPIDI_Accumulate(pami_context_t context, ++sync->started; - params.send.header.iov_base = &req->accum_headers[req->state.index]; + params.send.header.iov_base = &(((MPIDI_Win_MsgInfo *)req->accum_headers)[req->state.index]); params.send.data.iov_len = req->target.dt.map[req->state.index].DLOOP_VECTOR_LEN; params.send.data.iov_base = req->buffer + req->state.local_offset; diff --git a/src/mpid/pamid/src/onesided/mpid_win_get_accumulate.c b/src/mpid/pamid/src/onesided/mpid_win_get_accumulate.c new file mode 100644 index 0000000..085f3ce --- /dev/null +++ b/src/mpid/pamid/src/onesided/mpid_win_get_accumulate.c @@ -0,0 +1,517 @@ +/* begin_generated_IBM_copyright_prolog */ +/* */ +/* This is an automatically generated copyright prolog. */ +/* After initializing, DO NOT MODIFY OR MOVE */ +/* --------------------------------------------------------------- */ +/* Licensed Materials - Property of IBM */ +/* Blue Gene/Q 5765-PER 5765-PRP */ +/* */ +/* (C) Copyright IBM Corp. 2011, 2012 All Rights Reserved */ +/* US Government Users Restricted Rights - */ +/* Use, duplication, or disclosure restricted */ +/* by GSA ADP Schedule Contract with IBM Corp. */ +/* */ +/* --------------------------------------------------------------- */ +/* */ +/* end_generated_IBM_copyright_prolog */ +/* (C)Copyright IBM Corp. 2007, 2011 */ +/** + * \file src/onesided/mpid_win_accumulate.c + * \brief ??? + */ +#include "mpidi_onesided.h" + +static void +MPIDI_Win_GetAccSendAckDoneCB(pami_context_t context, + void * _buf, + pami_result_t result) +{ + MPIU_Free (_buf); +} + +static void +MPIDI_Win_GetAccumSendAck(pami_context_t context, + void * _info, + pami_result_t result) +{ + MPIDI_Win_GetAccMsgInfo *msginfo = (MPIDI_Win_GetAccMsgInfo *) _info; + pami_result_t rc = PAMI_SUCCESS; + + //Copy from msginfo->addr to a contiguous buffer + MPIDI_Datatype result_dt; + char *buffer = NULL; + + MPIDI_Win_datatype_basic(msginfo->result_count, + msginfo->result_datatype, + &result_dt); + + int use_map = 0; + buffer = MPIU_Malloc(result_dt.size); + if (result_dt.contig) + memcpy(buffer, (msginfo->addr + result_dt.true_lb), result_dt.size); + else + { + use_map = 1; + MPID_assert(buffer != NULL); + + int mpi_errno = 0; + mpi_errno = MPIR_Localcopy(msginfo->addr, + msginfo->count, + msginfo->type, + buffer, + result_dt.size, + MPI_CHAR); + MPID_assert(mpi_errno == MPI_SUCCESS); + } + + //Schedule sends to source to result buffer and trigger completion + //callback there + pami_send_t params = { + .send = { + .header = { + .iov_len = sizeof(MPIDI_Win_GetAccMsgInfo), + }, + .dispatch = MPIDI_Protocols_WinGetAccumAck, + .dest = msginfo->src_endpoint, + }, + .events = { + .cookie = buffer, + .local_fn = MPIDI_Win_GetAccSendAckDoneCB, //cleanup buffer + }, + }; + + int index = 0; + size_t local_offset = 0; + //Set the map + MPIDI_Win_datatype_map(&result_dt); + MPID_assert(result_dt.num_contig == msginfo->result_num_contig); + + while (index < result_dt.num_contig) { + params.send.header.iov_base = msginfo; + params.send.data.iov_len = result_dt.map[index].DLOOP_VECTOR_LEN; + params.send.data.iov_base = buffer + local_offset; + + rc = PAMI_Send(context, ¶ms); + MPID_assert(rc == PAMI_SUCCESS); + local_offset += params.send.data.iov_len; + ++index; + } + + /** Review PAMI_Send semantics and consider moving the free calls to + the completion callback*/ + //free msginfo + MPIU_Free(msginfo); + + if (use_map && result_dt.map != &result_dt.__map) + MPIU_Free (result_dt.map); +} + +void +MPIDI_WinGetAccumCB(pami_context_t context, + void * cookie, + const void * _msginfo, + size_t msginfo_size, + const void * sndbuf, + size_t sndlen, + pami_endpoint_t sender, + pami_recv_t * recv) +{ + MPID_assert(recv != NULL); + MPID_assert(sndbuf == NULL); + MPID_assert(msginfo_size == sizeof(MPIDI_Win_GetAccMsgInfo)); + MPID_assert(_msginfo != NULL); + MPIDI_Win_GetAccMsgInfo * msginfo = (MPIDI_Win_GetAccMsgInfo *) + MPIU_Malloc(sizeof(MPIDI_Win_GetAccMsgInfo)); + + *msginfo = *(MPIDI_Win_GetAccMsgInfo *)_msginfo; + msginfo->src_endpoint = sender; + + int null=0; + pami_type_t pami_type; + pami_data_function pami_op; + MPI_Op op = msginfo->op; + + MPIDI_Datatype_to_pami(msginfo->type, &pami_type, op, &pami_op, &null); + + recv->addr = msginfo->addr; + recv->type = pami_type; + recv->offset = 0; + recv->data_fn = pami_op; + recv->data_cookie = NULL; + recv->local_fn = NULL; + recv->cookie = NULL; + + if (msginfo->counter == 0) + //We will now allocate a tempbuf, copy local contents and start a + //send + MPIDI_Win_GetAccumSendAck (context, msginfo, PAMI_SUCCESS); + else + MPIU_Free(msginfo); +} + +static void +MPIDI_Win_GetAccDoneCB(pami_context_t context, + void * cookie, + pami_result_t result) +{ + MPIDI_Win_request *req = (MPIDI_Win_request*)cookie; + ++req->win->mpid.sync.complete; + ++req->origin.completed; + + if (req->origin.completed == + (req->result_num_contig + req->target.dt.num_contig)) + { + if (req->buffer_free) + MPIU_Free(req->buffer); + if (req->accum_headers) + MPIU_Free(req->accum_headers); + MPIU_Free (req); + } + MPIDI_Progress_signal(); +} + +void +MPIDI_WinGetAccumAckCB(pami_context_t context, + void * cookie, + const void * _msginfo, + size_t msginfo_size, + const void * sndbuf, + size_t sndlen, + pami_endpoint_t sender, + pami_recv_t * recv) +{ + MPID_assert(recv != NULL); + MPID_assert(sndbuf == NULL); + MPID_assert(msginfo_size == sizeof(MPIDI_Win_GetAccMsgInfo)); + MPID_assert(_msginfo != NULL); + const MPIDI_Win_GetAccMsgInfo * msginfo =(const MPIDI_Win_GetAccMsgInfo *)_msginfo; + + int null=0; + pami_type_t pami_type; + pami_data_function pami_op; + MPI_Op op = msginfo->op; + + MPIDI_Datatype_to_pami(msginfo->result_datatype, &pami_type, op, &pami_op, &null); + + recv->addr = msginfo->result_addr; + recv->type = pami_type; + recv->offset = 0; + recv->data_fn = PAMI_DATA_COPY; + recv->data_cookie = NULL; + recv->local_fn = MPIDI_Win_GetAccDoneCB; + recv->cookie = msginfo->req; +} + +static pami_result_t +MPIDI_Get_accumulate(pami_context_t context, + void * _req) +{ + MPIDI_Win_request *req = (MPIDI_Win_request*)_req; + pami_result_t rc; + void *map; + + pami_send_t params = { + .send = { + .header = { + .iov_len = sizeof(MPIDI_Win_GetAccMsgInfo), + }, + .dispatch = MPIDI_Protocols_WinGetAccum, + .dest = req->dest, + }, + .events = { + .cookie = req, + .remote_fn = MPIDI_Win_GetAccDoneCB, //When all accumulates have + //completed remotely + //complete accumulate + }, + }; + + struct MPIDI_Win_sync* sync = &req->win->mpid.sync; + TRACE_ERR("Start index=%u/%d l-addr=%p r-base=%p r-offset=%zu (sync->started=%u sync->complete=%u)\n", + req->state.index, req->target.dt.num_contig, req->buffer, req->win->mpid.info[req->target.rank].base_addr, req->offset, sync->started, sync->complete); + while (req->state.index < req->target.dt.num_contig) { + if (sync->started > sync->complete + MPIDI_Process.rma_pending) + { + TRACE_ERR("Bailing out; index=%u/%d sync->started=%u sync->complete=%u\n", + req->state.index, req->target.dt.num_contig, sync->started, sync->complete); + return PAMI_EAGAIN; + } + ++sync->started; + + params.send.header.iov_base = &(((MPIDI_Win_GetAccMsgInfo *)req->accum_headers)[req->state.index]); + params.send.data.iov_len = req->target.dt.map[req->state.index].DLOOP_VECTOR_LEN; + params.send.data.iov_base = req->buffer + req->state.local_offset; + + if (req->target.dt.num_contig - req->state.index == 1) { + map=NULL; + if (req->target.dt.map != &req->target.dt.__map) { + map=(void *) req->target.dt.map; + } + + rc = PAMI_Send(context, ¶ms); + MPID_assert(rc == PAMI_SUCCESS); + if (map) + MPIU_Free(map); + return PAMI_SUCCESS; + } else { + rc = PAMI_Send(context, ¶ms); + MPID_assert(rc == PAMI_SUCCESS); + req->state.local_offset += params.send.data.iov_len; + ++req->state.index; + } + } + + MPIDI_Win_datatype_unmap(&req->target.dt); + + return PAMI_SUCCESS; +} + + +/** + * \brief MPI-PAMI glue for MPI_GET_ACCUMULATE function + * + * According to the MPI Specification: + * + * Each datatype argument must be a predefined datatype or + * a derived datatype, where all basic components are of the + * same predefined datatype. Both datatype arguments must be + * constructed from the same predefined datatype. + * + * \param[in] origin_addr Source buffer + * \param[in] origin_count Number of datatype elements + * \param[in] origin_datatype Source datatype + * \param[in] target_rank Destination rank (target) + * \param[in] target_disp Displacement factor in target buffer + * \param[in] target_count Number of target datatype elements + * \param[in] target_datatype Destination datatype + * \param[in] op Operand to perform + * \param[in] win Window + * \return MPI_SUCCESS + */ +#undef FUNCNAME +#define FUNCNAME MPID_Get_accumulate +#undef FCNAME +#define FCNAME MPIU_QUOTE(FUNCNAME) +int +MPID_Get_accumulate(const void * origin_addr, + int origin_count, + MPI_Datatype origin_datatype, + void * result_addr, + int result_count, + MPI_Datatype result_datatype, + int target_rank, + MPI_Aint target_disp, + int target_count, + MPI_Datatype target_datatype, + MPI_Op op, + MPID_Win *win) +{ + int mpi_errno = MPI_SUCCESS; + + if (op == MPI_NO_OP) {//we just need to fetch data + mpi_errno = MPID_Get(result_addr, + result_count, + result_datatype, + target_rank, + target_disp, + target_count, + target_datatype, + win); + return mpi_errno; + } + + MPIDI_Win_request *req; + req = MPIU_Calloc0(1, MPIDI_Win_request); + req->win = win; + req->type = MPIDI_WIN_REQUEST_ACCUMULATE, + req->offset = target_disp*win->mpid.info[target_rank].disp_unit; +#ifdef __BGQ__ + /* PAMI limitation as it doesnt permit VA of 0 to be passed into + * memregion create, so we must pass base_va of heap computed from + * an SPI call instead. So the target offset must be adjusted */ + if (req->win->create_flavor == MPI_WIN_FLAVOR_DYNAMIC) + req->offset -= (size_t)req->win->mpid.info[target_rank].base_addr; +#endif + + if(win->mpid.sync.origin_epoch_type == win->mpid.sync.target_epoch_type && + win->mpid.sync.origin_epoch_type == MPID_EPOTYPE_REFENCE){ + win->mpid.sync.origin_epoch_type = MPID_EPOTYPE_FENCE; + win->mpid.sync.target_epoch_type = MPID_EPOTYPE_FENCE; + } + + if(win->mpid.sync.origin_epoch_type == MPID_EPOTYPE_NONE || + win->mpid.sync.origin_epoch_type == MPID_EPOTYPE_POST){ + MPIU_ERR_SETANDSTMT(mpi_errno, MPI_ERR_RMA_SYNC, + return mpi_errno, "**rmasync"); + } + + req->offset = target_disp * win->mpid.info[target_rank].disp_unit; + + if (origin_datatype == MPI_DOUBLE_INT) { + MPIDI_Win_datatype_basic(origin_count*2, + MPI_DOUBLE, + &req->origin.dt); + MPIDI_Win_datatype_basic(target_count*2, + MPI_DOUBLE, + &req->target.dt); + } + else if (origin_datatype == MPI_LONG_DOUBLE_INT) + { + MPIDI_Win_datatype_basic(origin_count*2, + MPI_LONG_DOUBLE, + &req->origin.dt); + MPIDI_Win_datatype_basic(target_count*2, + MPI_LONG_DOUBLE, + &req->target.dt); + } + else if (origin_datatype == MPI_LONG_INT) + { + MPIDI_Win_datatype_basic(origin_count*2, + MPI_LONG, + &req->origin.dt); + MPIDI_Win_datatype_basic(target_count*2, + MPI_LONG, + &req->target.dt); + } + else if (origin_datatype == MPI_SHORT_INT) + { + MPIDI_Win_datatype_basic(origin_count*2, + MPI_INT, + &req->origin.dt); + MPIDI_Win_datatype_basic(target_count*2, + MPI_INT, + &req->target.dt); + } + else + { + MPIDI_Win_datatype_basic(origin_count, + origin_datatype, + &req->origin.dt); + MPIDI_Win_datatype_basic(target_count, + target_datatype, + &req->target.dt); + } + + MPID_assert(req->origin.dt.size == req->target.dt.size); + + if ( (req->origin.dt.size == 0) || + (target_rank == MPI_PROC_NULL)) + { + MPIU_Free (req); + return MPI_SUCCESS; + } + + req->target.rank = target_rank; + + + if (req->origin.dt.contig) + { + req->buffer_free = 0; + req->buffer = (char*)origin_addr + req->origin.dt.true_lb; + } + else + { + req->buffer_free = 1; + req->buffer = MPIU_Malloc(req->origin.dt.size); + MPID_assert(req->buffer != NULL); + + int mpi_errno = 0; + mpi_errno = MPIR_Localcopy(origin_addr, + origin_count, + origin_datatype, + req->buffer, + req->origin.dt.size, + MPI_CHAR); + MPID_assert(mpi_errno == MPI_SUCCESS); + } + + pami_result_t rc; + pami_task_t task = MPID_VCR_GET_LPID(win->comm_ptr->vcr, target_rank); + if (win->mpid.sync.origin_epoch_type == MPID_EPOTYPE_START && + !MPIDI_valid_group_rank(task, win->mpid.sync.sc.group)) + { + MPIU_ERR_SETANDSTMT(mpi_errno, MPI_ERR_RMA_SYNC, + return mpi_errno, "**rmasync"); + } + + rc = PAMI_Endpoint_create(MPIDI_Client, task, 0, &req->dest); + MPID_assert(rc == PAMI_SUCCESS); + + + MPIDI_Win_datatype_map(&req->target.dt); + MPIDI_Datatype result_dt; + MPIDI_Win_datatype_basic(result_count, result_datatype, &result_dt); + req->result_num_contig = 1; + if (!result_dt.contig) + req->result_num_contig =result_dt.pointer->max_contig_blocks*result_count+1; + //We wait for #messages depending on target and result_datatype + win->mpid.sync.total += (req->result_num_contig + req->target.dt.num_contig); + + { + MPI_Datatype basic_type = MPI_DATATYPE_NULL; + MPID_Datatype_get_basic_type(origin_datatype, basic_type); + /* MPID_Datatype_get_basic_type() doesn't handle the struct types */ + if ((origin_datatype == MPI_FLOAT_INT) || + (origin_datatype == MPI_DOUBLE_INT) || + (origin_datatype == MPI_LONG_INT) || + (origin_datatype == MPI_SHORT_INT) || + (origin_datatype == MPI_LONG_DOUBLE_INT)) + { + MPID_assert(basic_type == MPI_DATATYPE_NULL); + basic_type = origin_datatype; + } + MPID_assert(basic_type != MPI_DATATYPE_NULL); + + MPI_Datatype result_basic_type = MPI_DATATYPE_NULL; + MPID_Datatype_get_basic_type(result_datatype, result_basic_type); + /* MPID_Datatype_get_basic_type() doesn't handle the struct types */ + if ((result_datatype == MPI_FLOAT_INT) || + (result_datatype == MPI_DOUBLE_INT) || + (result_datatype == MPI_LONG_INT) || + (result_datatype == MPI_SHORT_INT) || + (result_datatype == MPI_LONG_DOUBLE_INT)) + { + MPID_assert(result_basic_type == MPI_DATATYPE_NULL); + result_basic_type = result_datatype; + } + MPID_assert(result_basic_type != MPI_DATATYPE_NULL); + + unsigned index; + MPIDI_Win_GetAccMsgInfo * headers = MPIU_Calloc0(req->target.dt.num_contig, MPIDI_Win_GetAccMsgInfo); + req->accum_headers = headers; + for (index=0; index < req->target.dt.num_contig; ++index) { + headers[index].addr = win->mpid.info[target_rank].base_addr + req->offset + + (size_t)req->target.dt.map[index].DLOOP_VECTOR_BUF; + headers[index].req = req; + headers[index].win = win; + headers[index].type = basic_type; + headers[index].op = op; + headers[index].count = target_count; + headers[index].counter = index; + headers[index].num_contig = req->target.dt.num_contig; + headers[index].request = req; + headers[index].result_addr = result_addr; + headers[index].result_count = result_count; + headers[index].result_datatype = result_basic_type; + headers[index].result_num_contig= req->result_num_contig; + } + + } + + /* The pamid one-sided design requires context post in order to handle the + * case where the number of pending rma operation exceeds the + * 'PAMID_RMA_PENDING' threshold. When there are too many pending requests the + * work function remains on the context post queue (by returning PAMI_EAGAIN) + * so that the next time the context is advanced the work function will be + * invoked again. + * + * TODO - When context post is not required it would be better to attempt a + * direct context operation and then fail over to using context post if + * the rma pending threshold has been reached. This would result in + * better latency for one-sided operations. + */ + PAMI_Context_post(MPIDI_Context[0], &req->post_request, MPIDI_Get_accumulate, req); + +fn_fail: + return mpi_errno; +} diff --git a/src/mpid/pamid/src/onesided/mpidi_onesided.h b/src/mpid/pamid/src/onesided/mpidi_onesided.h index 8fa896b..e29b1e7 100644 --- a/src/mpid/pamid/src/onesided/mpidi_onesided.h +++ b/src/mpid/pamid/src/onesided/mpidi_onesided.h @@ -115,9 +115,27 @@ typedef struct size_t len; /* length of the send data */ } MPIDI_Win_MsgInfo; +typedef struct +{ + void * addr; + void * req; + MPID_Win * win; + MPI_Datatype type; + MPI_Op op; + int count; + int counter; + int num_contig; + void * request; + void * result_addr; + int result_count; + MPI_Datatype result_datatype; + int result_num_contig; + pami_endpoint_t src_endpoint; +} MPIDI_Win_GetAccMsgInfo; + /** \todo make sure that the extra fields are removed */ -typedef struct +typedef struct _mpidi_win_request { MPID_Win *win; MPIDI_Win_requesttype_t type; @@ -131,7 +149,7 @@ typedef struct size_t local_offset; } state; - MPIDI_Win_MsgInfo * accum_headers; + void * accum_headers; struct { @@ -153,9 +171,15 @@ typedef struct MPIDI_Datatype dt; } target; - void *buffer; void *user_buffer; uint32_t buffer_free; + void *buffer; + struct _mpidi_win_request *next; + void * compare_addr; + void * result_addr; + MPI_Op op; + int result_num_contig; + } MPIDI_Win_request; MPIDI_Win_request zero_req; /* used for init. request structure to 0 */ ----------------------------------------------------------------------- hooks/post-receive -- MPICH primary repository