#include "src/include/pmix_config.h"
#include "src/include/pmix_stdint.h"
#include "include/pmix.h"
#include "src/include/pmix_globals.h"
#ifdef HAVE_STRING_H
# include <string.h>
#endif
#include <fcntl.h>
#ifdef HAVE_UNISTD_H
# include <unistd.h>
#endif
#ifdef HAVE_SYS_SOCKET_H
# include <sys/socket.h>
#endif
#ifdef HAVE_SYS_UN_H
# include <sys/un.h>
#endif
#ifdef HAVE_SYS_UIO_H
# include <sys/uio.h>
#endif
#ifdef HAVE_SYS_TYPES_H
# include <sys/types.h>
#endif
#include <event.h>
#include "src/class/pmix_list.h"
#include "src/mca/bfrops/bfrops.h"
#include "src/mca/ptl/ptl.h"
#include "src/util/pmix_argv.h"
#include "src/util/pmix_error.h"
#include "src/util/pmix_output.h"
#include "pmix_client_ops.h"
static pmix_status_t unpack_return(pmix_buffer_t *data);
static pmix_status_t pack_fence(pmix_buffer_t *msg, pmix_cmd_t cmd, const pmix_proc_t *procs,
size_t nprocs, const pmix_info_t *info, size_t ninfo);
static void wait_cbfunc(struct pmix_peer_t *pr, pmix_ptl_hdr_t *hdr, pmix_buffer_t *buf,
void *cbdata);
static void op_cbfunc(pmix_status_t status, void *cbdata);
PMIX_EXPORT pmix_status_t PMIx_Fence(const pmix_proc_t procs[], size_t nprocs,
const pmix_info_t info[], size_t ninfo)
{
pmix_cb_t *cb;
pmix_status_t rc;
PMIX_ACQUIRE_THREAD(&pmix_global_lock);
pmix_output_verbose(2, pmix_client_globals.fence_output,
"pmix: executing fence");
if (pmix_globals.init_cntr <= 0) {
PMIX_RELEASE_THREAD(&pmix_global_lock);
return PMIX_ERR_INIT;
}
if (pmix_client_globals.singleton) {
PMIX_RELEASE_THREAD(&pmix_global_lock);
return PMIX_SUCCESS;
}
if (!pmix_globals.connected) {
PMIX_RELEASE_THREAD(&pmix_global_lock);
return PMIX_ERR_UNREACH;
}
PMIX_RELEASE_THREAD(&pmix_global_lock);
cb = PMIX_NEW(pmix_cb_t);
if (PMIX_SUCCESS != (rc = PMIx_Fence_nb(procs, nprocs, info, ninfo, op_cbfunc, cb))) {
PMIX_ERROR_LOG(rc);
PMIX_RELEASE(cb);
return rc;
}
PMIX_WAIT_THREAD(&cb->lock);
rc = cb->status;
PMIX_RELEASE(cb);
pmix_output_verbose(2, pmix_client_globals.fence_output,
"pmix: fence released");
return rc;
}
PMIX_EXPORT pmix_status_t PMIx_Fence_nb(const pmix_proc_t procs[], size_t nprocs,
const pmix_info_t info[], size_t ninfo,
pmix_op_cbfunc_t cbfunc, void *cbdata)
{
pmix_buffer_t *msg;
pmix_cmd_t cmd = PMIX_FENCENB_CMD;
pmix_status_t rc;
pmix_cb_t *cb;
pmix_proc_t rg, *rgs;
size_t nrg;
PMIX_ACQUIRE_THREAD(&pmix_global_lock);
pmix_output_verbose(2, pmix_client_globals.fence_output,
"pmix: fence_nb called");
if (pmix_globals.init_cntr <= 0) {
PMIX_RELEASE_THREAD(&pmix_global_lock);
return PMIX_ERR_INIT;
}
if (!pmix_globals.connected) {
PMIX_RELEASE_THREAD(&pmix_global_lock);
return PMIX_ERR_UNREACH;
}
PMIX_RELEASE_THREAD(&pmix_global_lock);
if (NULL == procs && 0 != nprocs) {
return PMIX_ERR_BAD_PARAM;
}
if (NULL == procs) {
pmix_strncpy(rg.nspace, pmix_globals.myid.nspace, PMIX_MAX_NSLEN);
rg.rank = PMIX_RANK_WILDCARD;
rgs = &rg;
nrg = 1;
} else {
rgs = (pmix_proc_t *) procs;
nrg = nprocs;
}
msg = PMIX_NEW(pmix_buffer_t);
if (PMIX_SUCCESS != (rc = pack_fence(msg, cmd, rgs, nrg, info, ninfo))) {
PMIX_RELEASE(msg);
return rc;
}
cb = PMIX_NEW(pmix_cb_t);
cb->cbfunc.opfn = cbfunc;
cb->cbdata = cbdata;
PMIX_PTL_SEND_RECV(rc, pmix_client_globals.myserver, msg, wait_cbfunc, (void *) cb);
if (PMIX_SUCCESS != rc) {
PMIX_RELEASE(msg);
PMIX_RELEASE(cb);
}
return rc;
}
static pmix_status_t unpack_return(pmix_buffer_t *data)
{
pmix_status_t rc;
pmix_status_t ret;
int32_t cnt;
pmix_output_verbose(2, pmix_client_globals.fence_output,
"client:unpack fence called");
cnt = 1;
PMIX_BFROPS_UNPACK(rc, pmix_client_globals.myserver, data, &ret, &cnt, PMIX_STATUS);
if (PMIX_SUCCESS != rc) {
PMIX_ERROR_LOG(rc);
return rc;
}
pmix_output_verbose(2, pmix_client_globals.fence_output,
"client:unpack fence received status %d", ret);
PMIX_GDS_RECV_MODEX_COMPLETE(rc, pmix_client_globals.myserver, data);
if (PMIX_SUCCESS != rc) {
PMIX_ERROR_LOG(rc);
return rc;
}
return ret;
}
static pmix_status_t pack_fence(pmix_buffer_t *msg, pmix_cmd_t cmd, const pmix_proc_t *procs,
size_t nprocs, const pmix_info_t *info, size_t ninfo)
{
pmix_status_t rc;
PMIX_BFROPS_PACK(rc, pmix_client_globals.myserver, msg, &cmd, 1, PMIX_COMMAND);
if (PMIX_SUCCESS != rc) {
PMIX_ERROR_LOG(rc);
return rc;
}
PMIX_BFROPS_PACK(rc, pmix_client_globals.myserver, msg, &nprocs, 1, PMIX_SIZE);
if (PMIX_SUCCESS != rc) {
PMIX_ERROR_LOG(rc);
return rc;
}
PMIX_BFROPS_PACK(rc, pmix_client_globals.myserver, msg, procs, nprocs, PMIX_PROC);
if (PMIX_SUCCESS != rc) {
PMIX_ERROR_LOG(rc);
return rc;
}
PMIX_BFROPS_PACK(rc, pmix_client_globals.myserver, msg, &ninfo, 1, PMIX_SIZE);
if (PMIX_SUCCESS != rc) {
PMIX_ERROR_LOG(rc);
return rc;
}
if (NULL != info && 0 < ninfo) {
PMIX_BFROPS_PACK(rc, pmix_client_globals.myserver, msg, info, ninfo, PMIX_INFO);
if (PMIX_SUCCESS != rc) {
PMIX_ERROR_LOG(rc);
return rc;
}
}
return PMIX_SUCCESS;
}
static void wait_cbfunc(struct pmix_peer_t *pr, pmix_ptl_hdr_t *hdr, pmix_buffer_t *buf,
void *cbdata)
{
pmix_cb_t *cb = (pmix_cb_t *) cbdata;
pmix_status_t rc;
PMIX_HIDE_UNUSED_PARAMS(pr, hdr);
pmix_output_verbose(2, pmix_client_globals.fence_output,
"pmix: fence_nb callback recvd");
if (NULL == cb) {
PMIX_ERROR_LOG(PMIX_ERR_BAD_PARAM);
return;
}
if (PMIX_BUFFER_IS_EMPTY(buf)) {
rc = PMIX_ERR_UNREACH;
} else {
rc = unpack_return(buf);
}
if (NULL != cb->cbfunc.opfn) {
cb->cbfunc.opfn(rc, cb->cbdata);
}
PMIX_RELEASE(cb);
}
static void op_cbfunc(pmix_status_t status, void *cbdata)
{
pmix_cb_t *cb = (pmix_cb_t *) cbdata;
cb->status = status;
PMIX_WAKEUP_THREAD(&cb->lock);
}