2015-06-18 19:53:20 +03:00
|
|
|
/*
|
|
|
|
* Copyright (c) 2004-2010 The Trustees of Indiana University and Indiana
|
|
|
|
* University Research and Technology
|
|
|
|
* Corporation. All rights reserved.
|
|
|
|
* Copyright (c) 2004-2011 The University of Tennessee and The University
|
|
|
|
* of Tennessee Research Foundation. All rights
|
|
|
|
* reserved.
|
|
|
|
* Copyright (c) 2004-2005 High Performance Computing Center Stuttgart,
|
|
|
|
* University of Stuttgart. All rights reserved.
|
|
|
|
* Copyright (c) 2004-2005 The Regents of the University of California.
|
|
|
|
* All rights reserved.
|
|
|
|
* Copyright (c) 2006-2013 Los Alamos National Security, LLC.
|
|
|
|
* All rights reserved.
|
|
|
|
* Copyright (c) 2009 Cisco Systems, Inc. All rights reserved.
|
|
|
|
* Copyright (c) 2011 Oak Ridge National Labs. All rights reserved.
|
2017-04-05 07:09:02 +03:00
|
|
|
* Copyright (c) 2013-2017 Intel, Inc. All rights reserved.
|
2015-06-18 19:53:20 +03:00
|
|
|
* Copyright (c) 2014 Mellanox Technologies, Inc.
|
|
|
|
* All rights reserved.
|
2016-03-28 09:07:01 +03:00
|
|
|
* Copyright (c) 2014-2016 Research Organization for Information Science
|
2015-06-18 19:53:20 +03:00
|
|
|
* and Technology (RIST). All rights reserved.
|
|
|
|
* $COPYRIGHT$
|
|
|
|
*
|
|
|
|
* Additional copyrights may follow
|
|
|
|
*
|
|
|
|
* $HEADER$
|
|
|
|
*
|
|
|
|
*/
|
|
|
|
|
|
|
|
#include "orte_config.h"
|
|
|
|
|
|
|
|
#ifdef HAVE_UNISTD_H
|
|
|
|
#include <unistd.h>
|
|
|
|
#endif
|
|
|
|
|
2015-12-16 02:26:13 +03:00
|
|
|
#include "opal/util/argv.h"
|
2015-06-18 19:53:20 +03:00
|
|
|
#include "opal/util/output.h"
|
|
|
|
#include "opal/dss/dss.h"
|
|
|
|
|
|
|
|
#include "orte/mca/errmgr/errmgr.h"
|
|
|
|
#include "orte/util/name_fns.h"
|
2017-04-05 14:27:32 +03:00
|
|
|
#include "orte/util/show_help.h"
|
2017-06-06 01:22:28 +03:00
|
|
|
#include "orte/util/threads.h"
|
2015-08-31 06:54:45 +03:00
|
|
|
#include "orte/runtime/orte_data_server.h"
|
2015-06-18 19:53:20 +03:00
|
|
|
#include "orte/runtime/orte_globals.h"
|
|
|
|
#include "orte/mca/rml/rml.h"
|
2017-05-26 18:57:55 +03:00
|
|
|
#include "orte/mca/rml/base/rml_contact.h"
|
2015-06-18 19:53:20 +03:00
|
|
|
|
|
|
|
#include "pmix_server_internal.h"
|
|
|
|
|
2017-05-26 18:57:55 +03:00
|
|
|
static int init_server(void)
|
|
|
|
{
|
|
|
|
char *server;
|
2017-07-06 19:48:48 +03:00
|
|
|
opal_value_t val;
|
2017-05-26 18:57:55 +03:00
|
|
|
char input[1024], *filename;
|
|
|
|
FILE *fp;
|
|
|
|
int rc;
|
|
|
|
|
|
|
|
/* only do this once */
|
|
|
|
orte_pmix_server_globals.pubsub_init = true;
|
|
|
|
|
|
|
|
/* if the universal server wasn't specified, then we use
|
|
|
|
* our own HNP for that purpose */
|
2017-05-27 20:47:08 +03:00
|
|
|
if (NULL == orte_data_server_uri) {
|
2017-05-26 18:57:55 +03:00
|
|
|
orte_pmix_server_globals.server = *ORTE_PROC_MY_HNP;
|
|
|
|
} else {
|
2017-05-27 20:47:08 +03:00
|
|
|
if (0 == strncmp(orte_data_server_uri, "file", strlen("file")) ||
|
|
|
|
0 == strncmp(orte_data_server_uri, "FILE", strlen("FILE"))) {
|
2017-05-26 18:57:55 +03:00
|
|
|
/* it is a file - get the filename */
|
2017-05-27 20:47:08 +03:00
|
|
|
filename = strchr(orte_data_server_uri, ':');
|
2017-05-26 18:57:55 +03:00
|
|
|
if (NULL == filename) {
|
|
|
|
/* filename is not correctly formatted */
|
|
|
|
orte_show_help("help-orterun.txt", "orterun:ompi-server-filename-bad", true,
|
2017-05-30 19:43:01 +03:00
|
|
|
orte_basename, orte_data_server_uri);
|
2017-05-26 18:57:55 +03:00
|
|
|
return ORTE_ERR_BAD_PARAM;
|
|
|
|
}
|
|
|
|
++filename; /* space past the : */
|
|
|
|
|
|
|
|
if (0 >= strlen(filename)) {
|
|
|
|
/* they forgot to give us the name! */
|
|
|
|
orte_show_help("help-orterun.txt", "orterun:ompi-server-filename-missing", true,
|
2017-05-30 19:43:01 +03:00
|
|
|
orte_basename, orte_data_server_uri);
|
2017-05-26 18:57:55 +03:00
|
|
|
return ORTE_ERR_BAD_PARAM;
|
|
|
|
}
|
|
|
|
|
|
|
|
/* open the file and extract the uri */
|
|
|
|
fp = fopen(filename, "r");
|
|
|
|
if (NULL == fp) { /* can't find or read file! */
|
|
|
|
orte_show_help("help-orterun.txt", "orterun:ompi-server-filename-access", true,
|
2017-05-30 19:43:01 +03:00
|
|
|
orte_basename, orte_data_server_uri);
|
2017-05-26 18:57:55 +03:00
|
|
|
return ORTE_ERR_BAD_PARAM;
|
|
|
|
}
|
|
|
|
if (NULL == fgets(input, 1024, fp)) {
|
|
|
|
/* something malformed about file */
|
|
|
|
fclose(fp);
|
|
|
|
orte_show_help("help-orterun.txt", "orterun:ompi-server-file-bad", true,
|
2017-05-30 19:43:01 +03:00
|
|
|
orte_basename, orte_data_server_uri,
|
2017-05-26 18:57:55 +03:00
|
|
|
orte_basename);
|
|
|
|
return ORTE_ERR_BAD_PARAM;
|
|
|
|
}
|
|
|
|
fclose(fp);
|
|
|
|
input[strlen(input)-1] = '\0'; /* remove newline */
|
|
|
|
server = strdup(input);
|
|
|
|
} else {
|
2017-05-30 19:43:01 +03:00
|
|
|
server = strdup(orte_data_server_uri);
|
2017-05-26 18:57:55 +03:00
|
|
|
}
|
2017-07-06 19:48:48 +03:00
|
|
|
/* parse the URI to get the server's name */
|
|
|
|
if (ORTE_SUCCESS != (rc = orte_rml_base_parse_uris(server, &orte_pmix_server_globals.server, NULL))) {
|
2017-05-26 18:57:55 +03:00
|
|
|
ORTE_ERROR_LOG(rc);
|
2017-07-22 01:33:16 +03:00
|
|
|
free(server);
|
2017-05-26 18:57:55 +03:00
|
|
|
return rc;
|
|
|
|
}
|
2017-07-06 19:48:48 +03:00
|
|
|
/* setup our route to the server */
|
|
|
|
OBJ_CONSTRUCT(&val, opal_value_t);
|
|
|
|
val.key = OPAL_PMIX_PROC_URI;
|
|
|
|
val.type = OPAL_STRING;
|
|
|
|
val.data.string = server;
|
|
|
|
if (OPAL_SUCCESS != (rc = opal_pmix.store_local(&orte_pmix_server_globals.server, &val))) {
|
2017-05-26 18:57:55 +03:00
|
|
|
ORTE_ERROR_LOG(rc);
|
2017-07-06 19:48:48 +03:00
|
|
|
val.key = NULL;
|
|
|
|
OBJ_DESTRUCT(&val);
|
2017-05-26 18:57:55 +03:00
|
|
|
return rc;
|
|
|
|
}
|
2017-07-06 19:48:48 +03:00
|
|
|
val.key = NULL;
|
|
|
|
OBJ_DESTRUCT(&val);
|
|
|
|
|
2017-05-26 18:57:55 +03:00
|
|
|
/* check if we are to wait for the server to start - resolves
|
|
|
|
* a race condition that can occur when the server is run
|
|
|
|
* as a background job - e.g., in scripts
|
|
|
|
*/
|
|
|
|
if (orte_pmix_server_globals.wait_for_server) {
|
|
|
|
/* ping the server */
|
|
|
|
struct timeval timeout;
|
|
|
|
timeout.tv_sec = orte_pmix_server_globals.timeout;
|
|
|
|
timeout.tv_usec = 0;
|
|
|
|
if (ORTE_SUCCESS != (rc = orte_rml.ping(orte_mgmt_conduit, server, &timeout))) {
|
|
|
|
/* try it one more time */
|
|
|
|
if (ORTE_SUCCESS != (rc = orte_rml.ping(orte_mgmt_conduit, server, &timeout))) {
|
|
|
|
/* okay give up */
|
|
|
|
orte_show_help("help-orterun.txt", "orterun:server-not-found", true,
|
|
|
|
orte_basename, server,
|
|
|
|
(long)orte_pmix_server_globals.timeout,
|
|
|
|
ORTE_ERROR_NAME(rc));
|
|
|
|
ORTE_UPDATE_EXIT_STATUS(ORTE_ERROR_DEFAULT_EXIT_CODE);
|
|
|
|
return rc;
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
return ORTE_SUCCESS;
|
|
|
|
}
|
|
|
|
|
2015-06-18 19:53:20 +03:00
|
|
|
static void execute(int sd, short args, void *cbdata)
|
|
|
|
{
|
|
|
|
pmix_server_req_t *req = (pmix_server_req_t*)cbdata;
|
|
|
|
int rc;
|
|
|
|
opal_buffer_t *xfer;
|
2017-05-26 18:57:55 +03:00
|
|
|
orte_process_name_t *target;
|
|
|
|
|
2017-06-06 01:22:28 +03:00
|
|
|
ORTE_ACQUIRE_OBJECT(req);
|
|
|
|
|
2017-05-26 18:57:55 +03:00
|
|
|
if (!orte_pmix_server_globals.pubsub_init) {
|
|
|
|
/* we need to initialize our connection to the server */
|
|
|
|
if (ORTE_SUCCESS != (rc = init_server())) {
|
|
|
|
orte_show_help("help-orted.txt", "noserver", true,
|
2017-05-30 19:43:01 +03:00
|
|
|
(NULL == orte_data_server_uri) ?
|
|
|
|
"NULL" : orte_data_server_uri);
|
2017-05-26 18:57:55 +03:00
|
|
|
goto callback;
|
|
|
|
}
|
|
|
|
}
|
2015-06-18 19:53:20 +03:00
|
|
|
|
|
|
|
/* add this request to our tracker hotel */
|
|
|
|
if (OPAL_SUCCESS != (rc = opal_hotel_checkin(&orte_pmix_server_globals.reqs, req, &req->room_num))) {
|
2017-04-05 14:27:32 +03:00
|
|
|
orte_show_help("help-orted.txt", "noroom", true, req->operation, orte_pmix_server_globals.num_rooms);
|
2015-06-18 19:53:20 +03:00
|
|
|
goto callback;
|
|
|
|
}
|
|
|
|
|
|
|
|
/* setup the xfer */
|
|
|
|
xfer = OBJ_NEW(opal_buffer_t);
|
|
|
|
/* pack the room number */
|
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(xfer, &req->room_num, 1, OPAL_INT))) {
|
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(xfer);
|
|
|
|
goto callback;
|
|
|
|
}
|
|
|
|
opal_dss.copy_payload(xfer, &req->msg);
|
|
|
|
|
2017-05-26 18:57:55 +03:00
|
|
|
/* if the range is SESSION, then set the target to the global server */
|
|
|
|
if (OPAL_PMIX_RANGE_SESSION == req->range) {
|
2017-09-19 03:41:27 +03:00
|
|
|
opal_output_verbose(1, orte_pmix_server_globals.output,
|
|
|
|
"%s orted:pmix:server range SESSION",
|
|
|
|
ORTE_NAME_PRINT(ORTE_PROC_MY_NAME));
|
2017-05-26 18:57:55 +03:00
|
|
|
target = &orte_pmix_server_globals.server;
|
2017-09-19 03:41:27 +03:00
|
|
|
} else if (OPAL_PMIX_RANGE_LOCAL == req->range) {
|
|
|
|
/* if the range is local, send it to myself */
|
|
|
|
opal_output_verbose(1, orte_pmix_server_globals.output,
|
|
|
|
"%s orted:pmix:server range LOCAL",
|
|
|
|
ORTE_NAME_PRINT(ORTE_PROC_MY_NAME));
|
|
|
|
target = ORTE_PROC_MY_NAME;
|
2017-05-26 18:57:55 +03:00
|
|
|
} else {
|
2017-09-19 03:41:27 +03:00
|
|
|
opal_output_verbose(1, orte_pmix_server_globals.output,
|
|
|
|
"%s orted:pmix:server range GLOBAL",
|
|
|
|
ORTE_NAME_PRINT(ORTE_PROC_MY_NAME));
|
2017-05-26 18:57:55 +03:00
|
|
|
target = ORTE_PROC_MY_HNP;
|
|
|
|
}
|
|
|
|
|
2015-06-18 19:53:20 +03:00
|
|
|
/* send the request to the target */
|
2016-10-23 21:36:05 +03:00
|
|
|
rc = orte_rml.send_buffer_nb(orte_mgmt_conduit,
|
2017-05-26 18:57:55 +03:00
|
|
|
target, xfer,
|
2015-06-18 19:53:20 +03:00
|
|
|
ORTE_RML_TAG_DATA_SERVER,
|
|
|
|
orte_rml_send_callback, NULL);
|
|
|
|
if (ORTE_SUCCESS == rc) {
|
|
|
|
return;
|
|
|
|
}
|
|
|
|
|
|
|
|
callback:
|
|
|
|
/* execute the callback to avoid having the client hang */
|
|
|
|
if (NULL != req->opcbfunc) {
|
|
|
|
req->opcbfunc(rc, req->cbdata);
|
|
|
|
} else if (NULL != req->lkcbfunc) {
|
|
|
|
req->lkcbfunc(rc, NULL, req->cbdata);
|
|
|
|
}
|
|
|
|
opal_hotel_checkout(&orte_pmix_server_globals.reqs, req->room_num);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
}
|
|
|
|
|
|
|
|
int pmix_server_publish_fn(opal_process_name_t *proc,
|
|
|
|
opal_list_t *info,
|
|
|
|
opal_pmix_op_cbfunc_t cbfunc, void *cbdata)
|
|
|
|
{
|
|
|
|
pmix_server_req_t *req;
|
|
|
|
int rc;
|
|
|
|
uint8_t cmd = ORTE_PMIX_PUBLISH_CMD;
|
2015-08-31 06:54:45 +03:00
|
|
|
opal_value_t *iptr;
|
2015-09-04 18:29:09 +03:00
|
|
|
opal_pmix_persistence_t persist = OPAL_PMIX_PERSIST_APP;
|
|
|
|
bool rset, pset;
|
2015-06-18 19:53:20 +03:00
|
|
|
|
2017-05-13 02:16:47 +03:00
|
|
|
opal_output_verbose(1, orte_pmix_server_globals.output,
|
|
|
|
"%s orted:pmix:server PUBLISH",
|
|
|
|
ORTE_NAME_PRINT(ORTE_PROC_MY_NAME));
|
|
|
|
|
2015-06-18 19:53:20 +03:00
|
|
|
/* create the caddy */
|
|
|
|
req = OBJ_NEW(pmix_server_req_t);
|
2017-04-05 07:09:02 +03:00
|
|
|
(void)asprintf(&req->operation, "PUBLISH: %s:%d", __FILE__, __LINE__);
|
2015-08-31 06:54:45 +03:00
|
|
|
req->opcbfunc = cbfunc;
|
|
|
|
req->cbdata = cbdata;
|
2015-06-18 19:53:20 +03:00
|
|
|
|
|
|
|
/* load the command */
|
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(&req->msg, &cmd, 1, OPAL_UINT8))) {
|
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
return rc;
|
|
|
|
}
|
|
|
|
|
|
|
|
/* pack the name of the publisher */
|
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(&req->msg, proc, 1, OPAL_NAME))) {
|
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
return rc;
|
|
|
|
}
|
|
|
|
|
2015-09-04 18:29:09 +03:00
|
|
|
/* no help for it - need to search for range/persistence */
|
|
|
|
rset = false;
|
|
|
|
pset = false;
|
|
|
|
OPAL_LIST_FOREACH(iptr, info, opal_value_t) {
|
|
|
|
if (0 == strcmp(iptr->key, OPAL_PMIX_RANGE)) {
|
2017-05-26 18:57:55 +03:00
|
|
|
req->range = (opal_pmix_data_range_t)iptr->data.uint;
|
2015-09-04 18:29:09 +03:00
|
|
|
if (pset) {
|
|
|
|
break;
|
|
|
|
}
|
|
|
|
rset = true;
|
|
|
|
} else if (0 == strcmp(iptr->key, OPAL_PMIX_PERSISTENCE)) {
|
2016-03-28 09:07:01 +03:00
|
|
|
persist = (opal_pmix_persistence_t)iptr->data.integer;
|
2015-09-04 18:29:09 +03:00
|
|
|
if (rset) {
|
|
|
|
break;
|
|
|
|
}
|
|
|
|
pset = true;
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2015-06-18 19:53:20 +03:00
|
|
|
/* pack the range */
|
2017-05-26 18:57:55 +03:00
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(&req->msg, &req->range, 1, OPAL_PMIX_DATA_RANGE))) {
|
2015-06-18 19:53:20 +03:00
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
return rc;
|
|
|
|
}
|
|
|
|
|
|
|
|
/* pack the persistence */
|
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(&req->msg, &persist, 1, OPAL_INT))) {
|
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
return rc;
|
|
|
|
}
|
|
|
|
|
2016-07-29 00:07:35 +03:00
|
|
|
/* if we have items, pack those too - ignore persistence, timeout
|
2015-09-04 18:29:09 +03:00
|
|
|
* and range values */
|
2015-08-31 06:54:45 +03:00
|
|
|
OPAL_LIST_FOREACH(iptr, info, opal_value_t) {
|
2015-09-04 18:29:09 +03:00
|
|
|
if (0 == strcmp(iptr->key, OPAL_PMIX_RANGE) ||
|
|
|
|
0 == strcmp(iptr->key, OPAL_PMIX_PERSISTENCE)) {
|
|
|
|
continue;
|
|
|
|
}
|
2016-07-29 00:07:35 +03:00
|
|
|
if (0 == strcmp(iptr->key, OPAL_PMIX_TIMEOUT)) {
|
|
|
|
/* record the timeout value, but don't pack it */
|
|
|
|
req->timeout = iptr->data.integer;
|
|
|
|
continue;
|
|
|
|
}
|
2015-09-13 22:59:26 +03:00
|
|
|
opal_output_verbose(5, orte_pmix_server_globals.output,
|
|
|
|
"%s publishing data %s of type %d from source %s",
|
|
|
|
ORTE_NAME_PRINT(ORTE_PROC_MY_NAME), iptr->key, iptr->type,
|
|
|
|
ORTE_NAME_PRINT(proc));
|
2015-08-31 06:54:45 +03:00
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(&req->msg, &iptr, 1, OPAL_VALUE))) {
|
2015-06-18 19:53:20 +03:00
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
return rc;
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
/* thread-shift so we can store the tracker */
|
|
|
|
opal_event_set(orte_event_base, &(req->ev),
|
|
|
|
-1, OPAL_EV_WRITE, execute, req);
|
|
|
|
opal_event_set_priority(&(req->ev), ORTE_MSG_PRI);
|
2017-06-06 01:22:28 +03:00
|
|
|
ORTE_POST_OBJECT(req);
|
2015-06-18 19:53:20 +03:00
|
|
|
opal_event_active(&(req->ev), OPAL_EV_WRITE, 1);
|
|
|
|
|
|
|
|
return OPAL_SUCCESS;
|
|
|
|
|
|
|
|
}
|
|
|
|
|
2015-09-04 18:29:09 +03:00
|
|
|
int pmix_server_lookup_fn(opal_process_name_t *proc, char **keys,
|
|
|
|
opal_list_t *info,
|
2015-06-18 19:53:20 +03:00
|
|
|
opal_pmix_lookup_cbfunc_t cbfunc, void *cbdata)
|
|
|
|
{
|
|
|
|
pmix_server_req_t *req;
|
|
|
|
int rc;
|
|
|
|
uint8_t cmd = ORTE_PMIX_LOOKUP_CMD;
|
2015-08-31 06:54:45 +03:00
|
|
|
int32_t nkeys, i;
|
|
|
|
opal_value_t *iptr;
|
|
|
|
|
|
|
|
/* the list of info objects are directives for us - they include
|
|
|
|
* things like timeout constraints, so there is no reason to
|
|
|
|
* forward them to the server */
|
2015-06-18 19:53:20 +03:00
|
|
|
|
|
|
|
/* create the caddy */
|
|
|
|
req = OBJ_NEW(pmix_server_req_t);
|
2017-04-05 07:09:02 +03:00
|
|
|
(void)asprintf(&req->operation, "LOOKUP: %s:%d", __FILE__, __LINE__);
|
2015-06-18 19:53:20 +03:00
|
|
|
req->lkcbfunc = cbfunc;
|
|
|
|
req->cbdata = cbdata;
|
|
|
|
|
|
|
|
/* load the command */
|
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(&req->msg, &cmd, 1, OPAL_UINT8))) {
|
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
return rc;
|
|
|
|
}
|
|
|
|
|
2016-06-25 03:01:49 +03:00
|
|
|
/* pack the requesting process jobid */
|
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(&req->msg, &proc->jobid, 1, ORTE_JOBID))) {
|
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
return rc;
|
|
|
|
}
|
|
|
|
|
2015-09-04 18:29:09 +03:00
|
|
|
/* no help for it - need to search for range */
|
|
|
|
OPAL_LIST_FOREACH(iptr, info, opal_value_t) {
|
|
|
|
if (0 == strcmp(iptr->key, OPAL_PMIX_RANGE)) {
|
2017-05-26 18:57:55 +03:00
|
|
|
req->range = (opal_pmix_data_range_t)iptr->data.uint;
|
2015-09-04 18:29:09 +03:00
|
|
|
break;
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2015-06-18 19:53:20 +03:00
|
|
|
/* pack the range */
|
2017-05-26 18:57:55 +03:00
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(&req->msg, &req->range, 1, OPAL_PMIX_DATA_RANGE))) {
|
2015-06-18 19:53:20 +03:00
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
return rc;
|
|
|
|
}
|
|
|
|
|
|
|
|
/* pack the number of keys */
|
|
|
|
nkeys = opal_argv_count(keys);
|
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(&req->msg, &nkeys, 1, OPAL_UINT32))) {
|
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
return rc;
|
|
|
|
}
|
|
|
|
|
|
|
|
/* pack the keys too */
|
2015-08-31 06:54:45 +03:00
|
|
|
for (i=0; i < nkeys; i++) {
|
2017-05-13 02:16:47 +03:00
|
|
|
opal_output_verbose(5, orte_pmix_server_globals.output,
|
|
|
|
"%s lookup data %s for proc %s",
|
|
|
|
ORTE_NAME_PRINT(ORTE_PROC_MY_NAME), keys[i],
|
|
|
|
ORTE_NAME_PRINT(proc));
|
2015-08-31 06:54:45 +03:00
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(&req->msg, &keys[i], 1, OPAL_STRING))) {
|
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
return rc;
|
|
|
|
}
|
2015-06-18 19:53:20 +03:00
|
|
|
}
|
|
|
|
|
2016-07-29 00:07:35 +03:00
|
|
|
/* if we have items, pack those too - ignore range and timeout value */
|
2015-09-04 18:29:09 +03:00
|
|
|
OPAL_LIST_FOREACH(iptr, info, opal_value_t) {
|
|
|
|
if (0 == strcmp(iptr->key, OPAL_PMIX_RANGE)) {
|
|
|
|
continue;
|
|
|
|
}
|
2016-07-29 00:07:35 +03:00
|
|
|
if (0 == strcmp(iptr->key, OPAL_PMIX_TIMEOUT)) {
|
|
|
|
/* record the timeout value, but don't pack it */
|
|
|
|
req->timeout = iptr->data.integer;
|
|
|
|
continue;
|
|
|
|
}
|
2017-08-17 21:58:48 +03:00
|
|
|
opal_output_verbose(2, orte_pmix_server_globals.output,
|
|
|
|
"%s lookup directive %s for proc %s",
|
|
|
|
ORTE_NAME_PRINT(ORTE_PROC_MY_NAME), iptr->key,
|
|
|
|
ORTE_NAME_PRINT(proc));
|
2015-09-04 18:29:09 +03:00
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(&req->msg, &iptr, 1, OPAL_VALUE))) {
|
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
return rc;
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2015-06-18 19:53:20 +03:00
|
|
|
/* thread-shift so we can store the tracker */
|
|
|
|
opal_event_set(orte_event_base, &(req->ev),
|
|
|
|
-1, OPAL_EV_WRITE, execute, req);
|
|
|
|
opal_event_set_priority(&(req->ev), ORTE_MSG_PRI);
|
2017-06-06 01:22:28 +03:00
|
|
|
ORTE_POST_OBJECT(req);
|
2015-06-18 19:53:20 +03:00
|
|
|
opal_event_active(&(req->ev), OPAL_EV_WRITE, 1);
|
|
|
|
|
|
|
|
return OPAL_SUCCESS;
|
|
|
|
}
|
|
|
|
|
2015-09-04 18:29:09 +03:00
|
|
|
int pmix_server_unpublish_fn(opal_process_name_t *proc, char **keys,
|
|
|
|
opal_list_t *info,
|
2015-06-18 19:53:20 +03:00
|
|
|
opal_pmix_op_cbfunc_t cbfunc, void *cbdata)
|
|
|
|
{
|
|
|
|
pmix_server_req_t *req;
|
|
|
|
int rc;
|
|
|
|
uint8_t cmd = ORTE_PMIX_UNPUBLISH_CMD;
|
2015-09-04 18:29:09 +03:00
|
|
|
uint32_t nkeys, n;
|
2015-08-31 06:54:45 +03:00
|
|
|
opal_value_t *iptr;
|
2015-06-18 19:53:20 +03:00
|
|
|
|
|
|
|
/* create the caddy */
|
|
|
|
req = OBJ_NEW(pmix_server_req_t);
|
2017-04-05 07:09:02 +03:00
|
|
|
(void)asprintf(&req->operation, "UNPUBLISH: %s:%d", __FILE__, __LINE__);
|
2015-08-31 06:54:45 +03:00
|
|
|
req->opcbfunc = cbfunc;
|
|
|
|
req->cbdata = cbdata;
|
2015-06-18 19:53:20 +03:00
|
|
|
|
|
|
|
/* load the command */
|
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(&req->msg, &cmd, 1, OPAL_UINT8))) {
|
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
return rc;
|
|
|
|
}
|
|
|
|
|
|
|
|
/* pack the name of the publisher */
|
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(&req->msg, proc, 1, OPAL_NAME))) {
|
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
return rc;
|
|
|
|
}
|
|
|
|
|
2015-09-04 18:29:09 +03:00
|
|
|
/* no help for it - need to search for range */
|
|
|
|
OPAL_LIST_FOREACH(iptr, info, opal_value_t) {
|
|
|
|
if (0 == strcmp(iptr->key, OPAL_PMIX_RANGE)) {
|
2017-05-26 18:57:55 +03:00
|
|
|
req->range = (opal_pmix_data_range_t)iptr->data.integer;
|
2015-09-04 18:29:09 +03:00
|
|
|
break;
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2015-06-18 19:53:20 +03:00
|
|
|
/* pack the range */
|
2017-05-26 18:57:55 +03:00
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(&req->msg, &req->range, 1, OPAL_INT))) {
|
2015-06-18 19:53:20 +03:00
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
return rc;
|
|
|
|
}
|
|
|
|
|
2015-09-04 18:29:09 +03:00
|
|
|
/* pack the number of keys */
|
|
|
|
nkeys = opal_argv_count(keys);
|
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(&req->msg, &nkeys, 1, OPAL_UINT32))) {
|
2015-08-31 06:54:45 +03:00
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
return rc;
|
|
|
|
}
|
|
|
|
|
2015-09-04 18:29:09 +03:00
|
|
|
/* pack the keys too */
|
|
|
|
for (n=0; n < nkeys; n++) {
|
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(&req->msg, &keys[n], 1, OPAL_STRING))) {
|
2015-08-31 06:54:45 +03:00
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
return rc;
|
|
|
|
}
|
|
|
|
}
|
2015-06-18 19:53:20 +03:00
|
|
|
|
2016-07-29 00:07:35 +03:00
|
|
|
/* if we have items, pack those too - ignore range and timeout value */
|
2015-09-04 18:29:09 +03:00
|
|
|
OPAL_LIST_FOREACH(iptr, info, opal_value_t) {
|
|
|
|
if (0 == strcmp(iptr->key, OPAL_PMIX_RANGE)) {
|
|
|
|
continue;
|
|
|
|
}
|
2016-07-29 00:07:35 +03:00
|
|
|
if (0 == strcmp(iptr->key, OPAL_PMIX_TIMEOUT)) {
|
|
|
|
/* record the timeout value, but don't pack it */
|
|
|
|
req->timeout = iptr->data.integer;
|
|
|
|
continue;
|
|
|
|
}
|
2015-09-04 18:29:09 +03:00
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.pack(&req->msg, &iptr, 1, OPAL_VALUE))) {
|
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(req);
|
|
|
|
return rc;
|
|
|
|
}
|
2015-06-18 19:53:20 +03:00
|
|
|
}
|
|
|
|
|
|
|
|
/* thread-shift so we can store the tracker */
|
|
|
|
opal_event_set(orte_event_base, &(req->ev),
|
|
|
|
-1, OPAL_EV_WRITE, execute, req);
|
|
|
|
opal_event_set_priority(&(req->ev), ORTE_MSG_PRI);
|
2017-06-06 01:22:28 +03:00
|
|
|
ORTE_POST_OBJECT(req);
|
2015-06-18 19:53:20 +03:00
|
|
|
opal_event_active(&(req->ev), OPAL_EV_WRITE, 1);
|
|
|
|
|
|
|
|
return OPAL_SUCCESS;
|
|
|
|
}
|
|
|
|
|
|
|
|
void pmix_server_keyval_client(int status, orte_process_name_t* sender,
|
|
|
|
opal_buffer_t *buffer,
|
|
|
|
orte_rml_tag_t tg, void *cbdata)
|
|
|
|
{
|
2015-09-05 21:19:41 +03:00
|
|
|
int rc, ret, room_num = -1;
|
2015-08-31 06:54:45 +03:00
|
|
|
int32_t cnt;
|
2015-09-05 21:19:41 +03:00
|
|
|
pmix_server_req_t *req=NULL;
|
2015-08-31 06:54:45 +03:00
|
|
|
opal_list_t info;
|
|
|
|
opal_value_t *iptr;
|
|
|
|
opal_pmix_pdata_t *pdata;
|
|
|
|
opal_process_name_t source;
|
2015-06-18 19:53:20 +03:00
|
|
|
|
2015-09-05 21:19:41 +03:00
|
|
|
opal_output_verbose(1, orte_pmix_server_globals.output,
|
|
|
|
"%s recvd lookup data return",
|
|
|
|
ORTE_NAME_PRINT(ORTE_PROC_MY_NAME));
|
|
|
|
|
2015-08-31 06:54:45 +03:00
|
|
|
OBJ_CONSTRUCT(&info, opal_list_t);
|
2015-06-18 19:53:20 +03:00
|
|
|
/* unpack the room number of the request tracker */
|
|
|
|
cnt = 1;
|
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.unpack(buffer, &room_num, &cnt, OPAL_INT))) {
|
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
return;
|
|
|
|
}
|
|
|
|
|
|
|
|
/* unpack the status */
|
|
|
|
cnt = 1;
|
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.unpack(buffer, &ret, &cnt, OPAL_INT))) {
|
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
ret = rc;
|
|
|
|
goto release;
|
|
|
|
}
|
|
|
|
|
2015-09-05 21:19:41 +03:00
|
|
|
opal_output_verbose(5, orte_pmix_server_globals.output,
|
|
|
|
"%s recvd lookup returned status %d",
|
|
|
|
ORTE_NAME_PRINT(ORTE_PROC_MY_NAME), ret);
|
|
|
|
|
2015-08-31 06:54:45 +03:00
|
|
|
if (ORTE_SUCCESS == ret) {
|
|
|
|
/* see if any data was included - not an error if the answer is no */
|
|
|
|
cnt = 1;
|
|
|
|
while (OPAL_SUCCESS == opal_dss.unpack(buffer, &source, &cnt, OPAL_NAME)) {
|
|
|
|
pdata = OBJ_NEW(opal_pmix_pdata_t);
|
|
|
|
pdata->proc = source;
|
|
|
|
if (OPAL_SUCCESS != (rc = opal_dss.unpack(buffer, &iptr, &cnt, OPAL_VALUE))) {
|
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(pdata);
|
|
|
|
continue;
|
2015-06-18 19:53:20 +03:00
|
|
|
}
|
2015-09-13 22:59:26 +03:00
|
|
|
opal_output_verbose(5, orte_pmix_server_globals.output,
|
|
|
|
"%s recvd lookup returned data %s of type %d from source %s",
|
|
|
|
ORTE_NAME_PRINT(ORTE_PROC_MY_NAME), iptr->key, iptr->type,
|
|
|
|
ORTE_NAME_PRINT(&source));
|
2015-08-31 06:54:45 +03:00
|
|
|
if (OPAL_SUCCESS != (rc = opal_value_xfer(&pdata->value, iptr))) {
|
|
|
|
ORTE_ERROR_LOG(rc);
|
|
|
|
OBJ_RELEASE(pdata);
|
|
|
|
OBJ_RELEASE(iptr);
|
|
|
|
continue;
|
2015-06-18 19:53:20 +03:00
|
|
|
}
|
2015-08-31 06:54:45 +03:00
|
|
|
OBJ_RELEASE(iptr);
|
|
|
|
opal_list_append(&info, &pdata->super);
|
2015-06-18 19:53:20 +03:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
release:
|
2015-09-05 21:19:41 +03:00
|
|
|
if (0 <= room_num) {
|
|
|
|
/* retrieve the tracker */
|
|
|
|
opal_hotel_checkout_and_return_occupant(&orte_pmix_server_globals.reqs, room_num, (void**)&req);
|
|
|
|
}
|
|
|
|
|
2015-08-30 07:19:27 +03:00
|
|
|
if (NULL != req) {
|
|
|
|
/* pass down the response */
|
|
|
|
if (NULL != req->opcbfunc) {
|
|
|
|
req->opcbfunc(ret, req->cbdata);
|
2015-08-31 06:54:45 +03:00
|
|
|
} else if (NULL != req->lkcbfunc) {
|
|
|
|
req->lkcbfunc(ret, &info, req->cbdata);
|
2015-08-30 07:19:27 +03:00
|
|
|
} else {
|
2015-08-31 06:54:45 +03:00
|
|
|
/* should not happen */
|
|
|
|
ORTE_ERROR_LOG(ORTE_ERR_NOT_SUPPORTED);
|
2015-08-30 07:19:27 +03:00
|
|
|
}
|
2015-06-18 19:53:20 +03:00
|
|
|
|
2015-08-30 07:19:27 +03:00
|
|
|
/* cleanup */
|
2015-08-31 06:54:45 +03:00
|
|
|
OPAL_LIST_DESTRUCT(&info);
|
2015-08-30 07:19:27 +03:00
|
|
|
OBJ_RELEASE(req);
|
|
|
|
}
|
2015-06-18 19:53:20 +03:00
|
|
|
}
|