#include "src/include/pmix_config.h"
#include "pmix_common.h"
#include "src/include/pmix_globals.h"
#include "src/class/pmix_list.h"
#include "src/util/pmix_argv.h"
#include "src/util/pmix_error.h"
#include "src/mca/gds/base/base.h"
#include "src/server/pmix_server_ops.h"
char *pmix_gds_base_get_available_modules(void)
{
if (!pmix_gds_globals.initialized) {
return NULL;
}
return strdup(pmix_gds_globals.all_mods);
}
pmix_gds_base_module_t *pmix_gds_base_assign_module(pmix_info_t *info, size_t ninfo)
{
pmix_gds_base_active_module_t *active;
pmix_gds_base_module_t *mod = NULL;
int pri, priority = -1;
if (!pmix_gds_globals.initialized) {
return NULL;
}
PMIX_LIST_FOREACH (active, &pmix_gds_globals.actives, pmix_gds_base_active_module_t) {
if (NULL == active->module->assign_module) {
continue;
}
if (PMIX_SUCCESS == active->module->assign_module(info, ninfo, &pri)) {
if (pri < 0) {
pri = active->pri;
}
if (priority < pri) {
mod = active->module;
priority = pri;
}
}
}
return mod;
}
pmix_status_t pmix_gds_base_setup_fork(const pmix_proc_t *proc, char ***env)
{
pmix_gds_base_active_module_t *active;
pmix_status_t rc;
if (!pmix_gds_globals.initialized) {
return PMIX_ERR_INIT;
}
PMIX_LIST_FOREACH (active, &pmix_gds_globals.actives, pmix_gds_base_active_module_t) {
if (NULL == active->module->setup_fork) {
continue;
}
rc = active->module->setup_fork(proc, env);
if (PMIX_SUCCESS != rc && PMIX_ERR_NOT_AVAILABLE != rc) {
return rc;
}
}
return PMIX_SUCCESS;
}
pmix_status_t pmix_gds_base_store_modex(struct pmix_namespace_t *nspace, pmix_buffer_t *buff,
pmix_gds_base_ctx_t ctx,
pmix_gds_base_store_modex_cb_fn_t cb_fn, void *cbdata)
{
(void) nspace;
pmix_status_t rc = PMIX_SUCCESS;
pmix_buffer_t bkt;
pmix_byte_object_t bo, bo2;
int32_t cnt = 1;
pmix_collect_t ctype;
pmix_server_trkr_t *trk = (pmix_server_trkr_t *) cbdata;
pmix_proc_t proc;
pmix_buffer_t pbkt;
pmix_rank_t rel_rank;
pmix_nspace_caddy_t *nm;
bool found;
char **kmap = NULL;
uint32_t kmap_size;
pmix_gds_modex_key_fmt_t kmap_type;
pmix_gds_modex_blob_info_t blob_info_byte = 0;
cnt = 1;
PMIX_BFROPS_UNPACK(rc, pmix_globals.mypeer, buff, &bo, &cnt, PMIX_BYTE_OBJECT);
if ((PMIX_COLLECT_YES == trk->collect_type) &&
(PMIX_ERR_UNPACK_READ_PAST_END_OF_BUFFER == rc)) {
goto exit;
}
while (PMIX_SUCCESS == rc) {
PMIX_CONSTRUCT(&bkt, pmix_buffer_t);
PMIX_LOAD_BUFFER(pmix_globals.mypeer, &bkt, bo.bytes, bo.size);
cnt = 1;
PMIX_BFROPS_UNPACK(rc, pmix_globals.mypeer, &bkt, &blob_info_byte, &cnt, PMIX_BYTE);
if (PMIX_ERR_UNPACK_READ_PAST_END_OF_BUFFER == rc) {
PMIX_DESTRUCT(&bkt);
break;
}
if (PMIX_SUCCESS != rc) {
PMIX_ERROR_LOG(rc);
PMIX_DESTRUCT(&bkt);
goto exit;
}
ctype = PMIX_GDS_COLLECT_IS_SET(blob_info_byte) ? PMIX_COLLECT_YES : PMIX_COLLECT_NO;
if (trk->collect_type != ctype) {
rc = PMIX_ERR_INVALID_ARG;
PMIX_ERROR_LOG(rc);
goto exit;
}
kmap_type = PMIX_GDS_KEYMAP_IS_SET(blob_info_byte) ? PMIX_MODEX_KEY_KEYMAP_FMT
: PMIX_MODEX_KEY_NATIVE_FMT;
if (PMIX_MODEX_KEY_KEYMAP_FMT == kmap_type) {
cnt = 1;
PMIX_BFROPS_UNPACK(rc, pmix_globals.mypeer, &bkt, &kmap_size, &cnt, PMIX_UINT32);
if (PMIX_ERR_UNPACK_READ_PAST_END_OF_BUFFER == rc) {
rc = PMIX_SUCCESS;
PMIX_DESTRUCT(&bkt);
break;
} else if (PMIX_SUCCESS != rc) {
PMIX_ERROR_LOG(rc);
PMIX_DESTRUCT(&bkt);
break;
}
kmap = (char **) (calloc(kmap_size + 1, sizeof(char *)));
if (NULL == kmap) {
rc = PMIX_ERR_OUT_OF_RESOURCE;
PMIX_ERROR_LOG(rc);
goto exit;
}
cnt = kmap_size;
PMIX_BFROPS_UNPACK(rc, pmix_globals.mypeer, &bkt, kmap, &cnt, PMIX_STRING);
if (PMIX_SUCCESS != rc) {
PMIX_ERROR_LOG(rc);
PMIX_DESTRUCT(&bkt);
goto exit;
}
if (PMIx_Argv_count(kmap) != (int) kmap_size) {
rc = PMIX_ERR_UNPACK_FAILURE;
PMIX_ERROR_LOG(rc);
PMIX_DESTRUCT(&bkt);
goto exit;
}
}
cnt = 1;
PMIX_BFROPS_UNPACK(rc, pmix_globals.mypeer, &bkt, &bo2, &cnt, PMIX_BYTE_OBJECT);
while (PMIX_SUCCESS == rc) {
PMIX_CONSTRUCT(&pbkt, pmix_buffer_t);
PMIX_LOAD_BUFFER(pmix_globals.mypeer, &pbkt, bo2.bytes, bo2.size);
cnt = 1;
PMIX_BFROPS_UNPACK(rc, pmix_globals.mypeer, &pbkt, &rel_rank, &cnt, PMIX_PROC_RANK);
if (PMIX_SUCCESS != rc) {
if (PMIX_ERR_UNPACK_READ_PAST_END_OF_BUFFER == rc) {
break;
}
PMIX_ERROR_LOG(rc);
pbkt.base_ptr = NULL;
PMIX_DESTRUCT(&pbkt);
break;
}
found = false;
if (pmix_list_get_size(&trk->nslist) == 1) {
found = true;
nm = (pmix_nspace_caddy_t *) pmix_list_get_first(&trk->nslist);
} else {
PMIX_LIST_FOREACH (nm, &trk->nslist, pmix_nspace_caddy_t) {
if (rel_rank < nm->ns->nprocs) {
found = true;
break;
}
rel_rank -= nm->ns->nprocs;
}
}
if (false == found) {
rc = PMIX_ERR_NOT_FOUND;
PMIX_ERROR_LOG(rc);
pbkt.base_ptr = NULL;
PMIX_DESTRUCT(&pbkt);
break;
}
PMIX_PROC_LOAD(&proc, nm->ns->nspace, rel_rank);
rc = cb_fn(ctx, &proc, kmap_type, kmap, &pbkt);
if (PMIX_SUCCESS != rc) {
PMIX_ERROR_LOG(rc);
pbkt.base_ptr = NULL;
PMIX_DESTRUCT(&pbkt);
break;
}
pbkt.base_ptr = NULL;
PMIX_DESTRUCT(&pbkt);
PMIX_BYTE_OBJECT_DESTRUCT(&bo2);
cnt = 1;
PMIX_BFROPS_UNPACK(rc, pmix_globals.mypeer, &bkt, &bo2, &cnt, PMIX_BYTE_OBJECT);
}
PMIX_DESTRUCT(&bkt);
if (PMIX_ERR_UNPACK_READ_PAST_END_OF_BUFFER == rc) {
rc = PMIX_SUCCESS;
} else if (PMIX_SUCCESS != rc) {
PMIX_ERROR_LOG(rc);
goto exit;
}
cnt = 1;
PMIX_BFROPS_UNPACK(rc, pmix_globals.mypeer, buff, &bo, &cnt, PMIX_BYTE_OBJECT);
}
if (PMIX_ERR_UNPACK_READ_PAST_END_OF_BUFFER == rc) {
rc = PMIX_SUCCESS;
} else if (PMIX_SUCCESS != rc) {
PMIX_ERROR_LOG(rc);
}
exit:
PMIx_Argv_free(kmap);
return rc;
}
pmix_status_t pmix_gds_base_modex_pack_kval(pmix_gds_modex_key_fmt_t key_fmt, pmix_buffer_t *buf,
char ***kmap, pmix_kval_t *kv)
{
uint32_t key_idx;
pmix_status_t rc = PMIX_SUCCESS;
if (PMIX_MODEX_KEY_KEYMAP_FMT == key_fmt) {
rc = pmix_argv_append_unique_idx((int *) &key_idx, kmap, kv->key);
if (PMIX_SUCCESS != rc) {
PMIX_ERROR_LOG(rc);
return rc;
}
PMIX_BFROPS_PACK(rc, pmix_globals.mypeer, buf, &key_idx, 1, PMIX_UINT32);
if (PMIX_SUCCESS != rc) {
PMIX_ERROR_LOG(rc);
return rc;
}
PMIX_BFROPS_PACK(rc, pmix_globals.mypeer, buf, kv->value, 1, PMIX_VALUE);
if (PMIX_SUCCESS != rc) {
PMIX_ERROR_LOG(rc);
return rc;
}
} else if (PMIX_MODEX_KEY_NATIVE_FMT == key_fmt) {
PMIX_BFROPS_PACK(rc, pmix_globals.mypeer, buf, kv, 1, PMIX_KVAL);
if (PMIX_SUCCESS != rc) {
PMIX_ERROR_LOG(rc);
return rc;
}
} else {
rc = PMIX_ERR_BAD_PARAM;
PMIX_ERROR_LOG(rc);
return rc;
}
return PMIX_SUCCESS;
}
pmix_status_t pmix_gds_base_modex_unpack_kval(pmix_gds_modex_key_fmt_t key_fmt, pmix_buffer_t *buf,
char **kmap, pmix_kval_t *kv)
{
int32_t cnt;
uint32_t key_idx;
pmix_status_t rc = PMIX_SUCCESS;
if (PMIX_MODEX_KEY_KEYMAP_FMT == key_fmt) {
cnt = 1;
PMIX_BFROPS_UNPACK(rc, pmix_globals.mypeer, buf, &key_idx, &cnt, PMIX_UINT32);
if (PMIX_SUCCESS != rc) {
return rc;
}
if (NULL == kmap[key_idx]) {
rc = PMIX_ERR_BAD_PARAM;
PMIX_ERROR_LOG(rc);
return rc;
}
kv->key = strdup(kmap[key_idx]);
cnt = 1;
PMIX_VALUE_CREATE(kv->value, 1);
PMIX_BFROPS_UNPACK(rc, pmix_globals.mypeer, buf, kv->value, &cnt, PMIX_VALUE);
if (PMIX_SUCCESS != rc) {
free(kv->key);
PMIX_VALUE_RELEASE(kv->value);
PMIX_ERROR_LOG(rc);
return rc;
}
} else if (PMIX_MODEX_KEY_NATIVE_FMT == key_fmt) {
cnt = 1;
PMIX_BFROPS_UNPACK(rc, pmix_globals.mypeer, buf, kv, &cnt, PMIX_KVAL);
if (PMIX_SUCCESS != rc) {
return rc;
}
} else {
rc = PMIX_ERR_BAD_PARAM;
PMIX_ERROR_LOG(rc);
return rc;
}
return PMIX_SUCCESS;
}