margo.c 37.3 KB
Newer Older
1 2 3 4 5 6 7 8 9 10 11

/*
 * (C) 2015 The University of Chicago
 * 
 * See COPYRIGHT in top-level directory.
 */

#include <assert.h>
#include <unistd.h>
#include <errno.h>
#include <abt.h>
Matthieu Dorier's avatar
Matthieu Dorier committed
12 13 14

#include <margo-config.h>
#ifdef HAVE_ABT_SNOOZER
15
#include <abt-snoozer.h>
Matthieu Dorier's avatar
Matthieu Dorier committed
16
#endif
17
#include <time.h>
Philip Carns's avatar
bug fix  
Philip Carns committed
18
#include <math.h>
19 20

#include "margo.h"
21
#include "margo-timer.h"
Philip Carns's avatar
Philip Carns committed
22
#include "utlist.h"
Philip Carns's avatar
Philip Carns committed
23
#include "uthash.h"
24

25
#define DEFAULT_MERCURY_PROGRESS_TIMEOUT_UB 100 /* 100 milliseconds */
Shane Snyder's avatar
Shane Snyder committed
26
#define DEFAULT_MERCURY_HANDLE_CACHE_SIZE 32
27

Philip Carns's avatar
Philip Carns committed
28 29 30 31 32 33 34 35 36 37
struct mplex_key
{
    hg_id_t id;
    uint32_t mplex_id;
};

struct mplex_element
{
    struct mplex_key key;
    ABT_pool pool;
38 39
    void* user_data;
    void(*user_free_callback)(void*);
Philip Carns's avatar
Philip Carns committed
40 41 42
    UT_hash_handle hh;
};

Philip Carns's avatar
Philip Carns committed
43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58
struct diag_data
{
    double min;
    double max;
    double cumulative;
    int count;
};

#define __DIAG_UPDATE(__data, __time)\
do {\
    __data.count++; \
    __data.cumulative += (__time); \
    if((__time) > __data.max) __data.max = (__time); \
    if((__time) < __data.min) __data.min = (__time); \
} while(0)

Shane Snyder's avatar
Shane Snyder committed
59 60 61 62 63 64 65
struct margo_handle_cache_el
{
    hg_handle_t handle;
    UT_hash_handle hh; /* in-use hash link */
    struct margo_handle_cache_el *next; /* free list link */
};

66 67 68 69 70 71 72
struct margo_finalize_cb
{
    void(*callback)(void*);
    void* uargs;
    struct margo_finalize_cb* next;
};

73 74
struct margo_timer_list; /* defined in margo-timer.c */

75 76
struct margo_instance
{
Shane Snyder's avatar
Shane Snyder committed
77
    /* mercury/argobots state */
78 79
    hg_context_t *hg_context;
    hg_class_t *hg_class;
80 81 82
    ABT_pool handler_pool;
    ABT_pool progress_pool;

83
    /* internal to margo for this particular instance */
Shane Snyder's avatar
Shane Snyder committed
84
    int margo_init;
85
    int abt_init;
86 87
    ABT_thread hg_progress_tid;
    int hg_progress_shutdown_flag;
Jonathan Jenkins's avatar
Jonathan Jenkins committed
88
    ABT_xstream progress_xstream;
Jonathan Jenkins's avatar
Jonathan Jenkins committed
89 90 91
    int owns_progress_pool;
    ABT_xstream *rpc_xstreams;
    int num_handler_pool_threads;
92
    unsigned int hg_progress_timeout_ub;
93 94 95

    /* control logic for callers waiting on margo to be finalized */
    int finalize_flag;
Jonathan Jenkins's avatar
Jonathan Jenkins committed
96
    int refcount;
97 98
    ABT_mutex finalize_mutex;
    ABT_cond finalize_cond;
99
    struct margo_finalize_cb* finalize_cb;
100

101 102 103
    /* timer data */
    struct margo_timer_list* timer_list;

Philip Carns's avatar
Philip Carns committed
104 105
    /* hash table to track multiplexed rpcs registered with margo */
    struct mplex_element *mplex_table;
Philip Carns's avatar
Philip Carns committed
106

Shane Snyder's avatar
Shane Snyder committed
107 108 109
    /* linked list of free hg handles and a hash of in-use handles */
    struct margo_handle_cache_el *free_handle_list;
    struct margo_handle_cache_el *used_handle_hash;
110
    ABT_mutex handle_cache_mtx; /* mutex protecting access to above caches */
Shane Snyder's avatar
Shane Snyder committed
111

Philip Carns's avatar
Philip Carns committed
112 113 114 115 116 117 118 119 120 121 122
    /* optional diagnostics data tracking */
    /* NOTE: technically the following fields are subject to races if they
     * are updated from more than one thread at a time.  We will be careful
     * to only update the counters from the progress_fn,
     * which will serialize access.
     */
    int diag_enabled;
    struct diag_data diag_trigger_elapsed;
    struct diag_data diag_progress_elapsed_zero_timeout;
    struct diag_data diag_progress_elapsed_nonzero_timeout;
    struct diag_data diag_progress_timeout_value;
123 124
};

125 126 127 128 129 130
struct margo_rpc_data
{
	margo_instance_id mid;
	void* user_data;
	void (*user_free_callback)(void *);
};
131

132
static void hg_progress_fn(void* foo);
133
static void margo_rpc_data_free(void* ptr);
134

Shane Snyder's avatar
Shane Snyder committed
135 136 137 138 139 140
static hg_return_t margo_handle_cache_init(margo_instance_id mid);
static void margo_handle_cache_destroy(margo_instance_id mid);
static hg_return_t margo_handle_cache_get(margo_instance_id mid,
    hg_addr_t addr, hg_id_t id, hg_handle_t *handle);
static hg_return_t margo_handle_cache_put(margo_instance_id mid,
    hg_handle_t handle);
141 142
static void delete_multiplexing_hash(margo_instance_id mid);

Shane Snyder's avatar
Shane Snyder committed
143

Shane Snyder's avatar
Shane Snyder committed
144
margo_instance_id margo_init(const char *addr_str, int mode,
Shane Snyder's avatar
Shane Snyder committed
145
    int use_progress_thread, int rpc_thread_count)
146
{
Jonathan Jenkins's avatar
Jonathan Jenkins committed
147 148 149 150 151
    ABT_xstream progress_xstream = ABT_XSTREAM_NULL;
    ABT_pool progress_pool = ABT_POOL_NULL;
    ABT_xstream *rpc_xstreams = NULL;
    ABT_xstream rpc_xstream = ABT_XSTREAM_NULL;
    ABT_pool rpc_pool = ABT_POOL_NULL;
Shane Snyder's avatar
Shane Snyder committed
152 153
    hg_class_t *hg_class = NULL;
    hg_context_t *hg_context = NULL;
Shane Snyder's avatar
Shane Snyder committed
154
    int listen_flag = (mode == MARGO_CLIENT_MODE) ? HG_FALSE : HG_TRUE;
155
    int abt_init = 0;
Jonathan Jenkins's avatar
Jonathan Jenkins committed
156
    int i;
Shane Snyder's avatar
Shane Snyder committed
157 158 159
    int ret;
    struct margo_instance *mid = MARGO_INSTANCE_NULL;

Shane Snyder's avatar
Shane Snyder committed
160
    if(mode != MARGO_CLIENT_MODE && mode != MARGO_SERVER_MODE) goto err;
Shane Snyder's avatar
Shane Snyder committed
161

162 163 164 165 166 167
    if (ABT_initialized() == ABT_ERR_UNINITIALIZED)
    {
        ret = ABT_init(0, NULL); /* XXX: argc/argv not currently used by ABT ... */
        if(ret != 0) goto err;
        abt_init = 1;
    }
Shane Snyder's avatar
Shane Snyder committed
168

169
    /* set caller (self) ES to idle without polling */
Matthieu Dorier's avatar
Matthieu Dorier committed
170
#ifdef HAVE_ABT_SNOOZER
Shane Snyder's avatar
Shane Snyder committed
171 172
    ret = ABT_snoozer_xstream_self_set();
    if(ret != 0) goto err;
Matthieu Dorier's avatar
Matthieu Dorier committed
173
#endif
Jonathan Jenkins's avatar
Jonathan Jenkins committed
174 175 176

    if (use_progress_thread)
    {
Matthieu Dorier's avatar
Matthieu Dorier committed
177
#ifdef HAVE_ABT_SNOOZER
Jonathan Jenkins's avatar
Jonathan Jenkins committed
178
        ret = ABT_snoozer_xstream_create(1, &progress_pool, &progress_xstream);
Matthieu Dorier's avatar
Matthieu Dorier committed
179 180 181 182 183 184 185
		if (ret != ABT_SUCCESS) goto err;
#else
		ret = ABT_xstream_create(ABT_SCHED_NULL, &progress_xstream);
		if (ret != ABT_SUCCESS) goto err;
		ret = ABT_xstream_get_main_pools(progress_xstream, 1, &progress_pool);
		if (ret != ABT_SUCCESS) goto err;
#endif
Jonathan Jenkins's avatar
Jonathan Jenkins committed
186 187 188 189 190 191 192 193 194
    }
    else
    {
        ret = ABT_xstream_self(&progress_xstream);
        if (ret != ABT_SUCCESS) goto err;
        ret = ABT_xstream_get_main_pools(progress_xstream, 1, &progress_pool);
        if (ret != ABT_SUCCESS) goto err;
    }

195
    if (rpc_thread_count > 0)
Jonathan Jenkins's avatar
Jonathan Jenkins committed
196
    {
197 198
        rpc_xstreams = calloc(rpc_thread_count, sizeof(*rpc_xstreams));
        if (rpc_xstreams == NULL) goto err;
Matthieu Dorier's avatar
Matthieu Dorier committed
199
#ifdef HAVE_ABT_SNOOZER
200 201 202
        ret = ABT_snoozer_xstream_create(rpc_thread_count, &rpc_pool,
                rpc_xstreams);
        if (ret != ABT_SUCCESS) goto err;
Matthieu Dorier's avatar
Matthieu Dorier committed
203
#else
204 205 206 207 208
        int j;
        ret = ABT_pool_create_basic(ABT_POOL_FIFO, ABT_POOL_ACCESS_MPMC, ABT_TRUE, &rpc_pool);
        if (ret != ABT_SUCCESS) goto err;
        for(j=0; j<rpc_thread_count; j++) {
            ret = ABT_xstream_create(ABT_SCHED_NULL, rpc_xstreams+j);
Shane Snyder's avatar
Shane Snyder committed
209 210
            if (ret != ABT_SUCCESS) goto err;
        }
211 212 213 214 215 216 217 218 219 220 221 222
#endif
    }
    else if (rpc_thread_count == 0)
    {
        ret = ABT_xstream_self(&rpc_xstream);
        if (ret != ABT_SUCCESS) goto err;
        ret = ABT_xstream_get_main_pools(rpc_xstream, 1, &rpc_pool);
        if (ret != ABT_SUCCESS) goto err;
    }
    else
    {
        rpc_pool = progress_pool;
Jonathan Jenkins's avatar
Jonathan Jenkins committed
223 224
    }

Shane Snyder's avatar
Shane Snyder committed
225 226 227 228 229 230
    hg_class = HG_Init(addr_str, listen_flag);
    if(!hg_class) goto err;

    hg_context = HG_Context_create(hg_class);
    if(!hg_context) goto err;

Jonathan Jenkins's avatar
Jonathan Jenkins committed
231 232 233
    mid = margo_init_pool(progress_pool, rpc_pool, hg_context);
    if (mid == MARGO_INSTANCE_NULL) goto err;

Shane Snyder's avatar
Shane Snyder committed
234
    mid->margo_init = 1;
235
    mid->abt_init = abt_init;
Jonathan Jenkins's avatar
Jonathan Jenkins committed
236 237 238
    mid->owns_progress_pool = use_progress_thread;
    mid->progress_xstream = progress_xstream;
    mid->num_handler_pool_threads = rpc_thread_count < 0 ? 0 : rpc_thread_count;
Jonathan Jenkins's avatar
Jonathan Jenkins committed
239
    mid->rpc_xstreams = rpc_xstreams;
240

Jonathan Jenkins's avatar
Jonathan Jenkins committed
241 242 243
    return mid;

err:
Shane Snyder's avatar
Shane Snyder committed
244 245
    if(mid)
    {
246
        margo_timer_list_free(mid->timer_list);
Shane Snyder's avatar
Shane Snyder committed
247 248 249 250
        ABT_mutex_free(&mid->finalize_mutex);
        ABT_cond_free(&mid->finalize_cond);
        free(mid);
    }
Jonathan Jenkins's avatar
Jonathan Jenkins committed
251 252 253 254 255 256 257 258 259 260 261 262 263 264
    if (use_progress_thread && progress_xstream != ABT_XSTREAM_NULL)
    {
        ABT_xstream_join(progress_xstream);
        ABT_xstream_free(&progress_xstream);
    }
    if (rpc_thread_count > 0 && rpc_xstreams != NULL)
    {
        for (i = 0; i < rpc_thread_count; i++)
        {
            ABT_xstream_join(rpc_xstreams[i]);
            ABT_xstream_free(&rpc_xstreams[i]);
        }
        free(rpc_xstreams);
    }
Shane Snyder's avatar
Shane Snyder committed
265 266 267 268
    if(hg_context)
        HG_Context_destroy(hg_context);
    if(hg_class)
        HG_Finalize(hg_class);
269 270
    if(abt_init)
        ABT_finalize();
Jonathan Jenkins's avatar
Jonathan Jenkins committed
271 272 273 274
    return MARGO_INSTANCE_NULL;
}

margo_instance_id margo_init_pool(ABT_pool progress_pool, ABT_pool handler_pool,
Jonathan Jenkins's avatar
Jonathan Jenkins committed
275
    hg_context_t *hg_context)
276 277
{
    int ret;
Shane Snyder's avatar
Shane Snyder committed
278
    hg_return_t hret;
279 280 281
    struct margo_instance *mid;

    mid = malloc(sizeof(*mid));
Shane Snyder's avatar
Shane Snyder committed
282
    if(!mid) goto err;
283
    memset(mid, 0, sizeof(*mid));
284

285 286 287
    ABT_mutex_create(&mid->finalize_mutex);
    ABT_cond_create(&mid->finalize_cond);

288 289
    mid->progress_pool = progress_pool;
    mid->handler_pool = handler_pool;
Jonathan Jenkins's avatar
Jonathan Jenkins committed
290
    mid->hg_class = HG_Context_get_class(hg_context);
291
    mid->hg_context = hg_context;
292
    mid->hg_progress_timeout_ub = DEFAULT_MERCURY_PROGRESS_TIMEOUT_UB;
293
    mid->mplex_table = NULL;
Jonathan Jenkins's avatar
Jonathan Jenkins committed
294
    mid->refcount = 1;
295
    mid->finalize_cb = NULL;
296

297 298
    mid->timer_list = margo_timer_list_create();
    if(mid->timer_list == NULL) goto err;
299

Shane Snyder's avatar
Shane Snyder committed
300 301 302 303
    /* initialize the handle cache */
    hret = margo_handle_cache_init(mid);
    if(hret != HG_SUCCESS) goto err;

304
    ret = ABT_thread_create(mid->progress_pool, hg_progress_fn, mid, 
305
        ABT_THREAD_ATTR_NULL, &mid->hg_progress_tid);
Shane Snyder's avatar
Shane Snyder committed
306 307
    if(ret != 0) goto err;

Shane Snyder's avatar
Shane Snyder committed
308 309
    return mid;

Shane Snyder's avatar
Shane Snyder committed
310 311
err:
    if(mid)
312
    {
Shane Snyder's avatar
Shane Snyder committed
313
        margo_handle_cache_destroy(mid);
314
        margo_timer_list_free(mid->timer_list);
Shane Snyder's avatar
Shane Snyder committed
315 316
        ABT_mutex_free(&mid->finalize_mutex);
        ABT_cond_free(&mid->finalize_cond);
317
        free(mid);
318
    }
Shane Snyder's avatar
Shane Snyder committed
319
    return MARGO_INSTANCE_NULL;
320 321
}

Jonathan Jenkins's avatar
Jonathan Jenkins committed
322 323 324 325
static void margo_cleanup(margo_instance_id mid)
{
    int i;

326 327 328 329 330 331 332 333 334
    /* call finalize callbacks */
    struct margo_finalize_cb* fcb = mid->finalize_cb;
    while(fcb) {
        (fcb->callback)(fcb->uargs);
        struct margo_finalize_cb* tmp = fcb;
        fcb = fcb->next;
        free(tmp);
    }

335
    margo_timer_list_free(mid->timer_list);
Jonathan Jenkins's avatar
Jonathan Jenkins committed
336

337 338 339
    /* delete the hash used for multiplexing */
    delete_multiplexing_hash(mid);

Jonathan Jenkins's avatar
Jonathan Jenkins committed
340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358
    ABT_mutex_free(&mid->finalize_mutex);
    ABT_cond_free(&mid->finalize_cond);

    if (mid->owns_progress_pool)
    {
        ABT_xstream_join(mid->progress_xstream);
        ABT_xstream_free(&mid->progress_xstream);
    }

    if (mid->num_handler_pool_threads > 0)
    {
        for (i = 0; i < mid->num_handler_pool_threads; i++)
        {
            ABT_xstream_join(mid->rpc_xstreams[i]);
            ABT_xstream_free(&mid->rpc_xstreams[i]);
        }
        free(mid->rpc_xstreams);
    }

Shane Snyder's avatar
Shane Snyder committed
359 360
    margo_handle_cache_destroy(mid);

Shane Snyder's avatar
Shane Snyder committed
361 362 363 364 365 366
    if (mid->margo_init)
    {
        if (mid->hg_context)
            HG_Context_destroy(mid->hg_context);
        if (mid->hg_class)
            HG_Finalize(mid->hg_class);
367 368
        if (mid->abt_init)
            ABT_finalize();
Shane Snyder's avatar
Shane Snyder committed
369 370
    }

Jonathan Jenkins's avatar
Jonathan Jenkins committed
371 372 373
    free(mid);
}

374
void margo_finalize(margo_instance_id mid)
375
{
Jonathan Jenkins's avatar
Jonathan Jenkins committed
376
    int do_cleanup;
377

378
    /* tell progress thread to wrap things up */
379
    mid->hg_progress_shutdown_flag = 1;
380 381

    /* wait for it to shutdown cleanly */
382 383
    ABT_thread_join(mid->hg_progress_tid);
    ABT_thread_free(&mid->hg_progress_tid);
384

385 386 387 388
    ABT_mutex_lock(mid->finalize_mutex);
    mid->finalize_flag = 1;
    ABT_cond_broadcast(mid->finalize_cond);

Jonathan Jenkins's avatar
Jonathan Jenkins committed
389 390
    mid->refcount--;
    do_cleanup = mid->refcount == 0;
391

Jonathan Jenkins's avatar
Jonathan Jenkins committed
392 393 394 395 396 397 398
    ABT_mutex_unlock(mid->finalize_mutex);

    /* if there was noone waiting on the finalize at the time of the finalize
     * broadcast, then we're safe to clean up. Otherwise, let the finalizer do
     * it */
    if (do_cleanup)
        margo_cleanup(mid);
399 400 401 402 403 404

    return;
}

void margo_wait_for_finalize(margo_instance_id mid)
{
Jonathan Jenkins's avatar
Jonathan Jenkins committed
405
    int do_cleanup;
406 407 408

    ABT_mutex_lock(mid->finalize_mutex);

Jonathan Jenkins's avatar
Jonathan Jenkins committed
409
        mid->refcount++;
410 411 412 413
            
        while(!mid->finalize_flag)
            ABT_cond_wait(mid->finalize_cond, mid->finalize_mutex);

Jonathan Jenkins's avatar
Jonathan Jenkins committed
414 415 416
        mid->refcount--;
        do_cleanup = mid->refcount == 0;

417
    ABT_mutex_unlock(mid->finalize_mutex);
Jonathan Jenkins's avatar
Jonathan Jenkins committed
418 419 420 421

    if (do_cleanup)
        margo_cleanup(mid);

422 423 424
    return;
}

425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441
void margo_push_finalize_callback(
            margo_instance_id mid,
            void(*cb)(void*),                  
            void* uargs)
{
    if(cb == NULL) return;

    struct margo_finalize_cb* fcb = 
        (struct margo_finalize_cb*)malloc(sizeof(*fcb));
    fcb->callback = cb;
    fcb->uargs = uargs;

    struct margo_finalize_cb* next = mid->finalize_cb;
    fcb->next = next;
    mid->finalize_cb = fcb;
}

442 443
hg_id_t margo_register_name(margo_instance_id mid, const char *func_name,
    hg_proc_cb_t in_proc_cb, hg_proc_cb_t out_proc_cb, hg_rpc_cb_t rpc_cb)
444
{
445 446 447
	struct margo_rpc_data* margo_data;
    hg_return_t hret;
    hg_id_t id;
448

449 450 451
    id = HG_Register_name(mid->hg_class, func_name, in_proc_cb, out_proc_cb, rpc_cb);
    if(id <= 0)
        return(0);
452

453 454 455 456 457 458 459 460 461 462 463 464
	/* register the margo data with the RPC */
    margo_data = (struct margo_rpc_data*)HG_Registered_data(mid->hg_class, id);
    if(!margo_data)
    {
        margo_data = (struct margo_rpc_data*)malloc(sizeof(struct margo_rpc_data));
        if(!margo_data)
            return(0);
        margo_data->mid = mid;
        margo_data->user_data = NULL;
        margo_data->user_free_callback = NULL;
        hret = HG_Register_data(mid->hg_class, id, margo_data, margo_rpc_data_free);
        if(hret != HG_SUCCESS)
465
        {
466 467
            free(margo_data);
            return(0);
468
        }
469 470
    }

471
	return(id);
472 473
}

474 475 476
hg_id_t margo_register_name_mplex(margo_instance_id mid, const char *func_name,
    hg_proc_cb_t in_proc_cb, hg_proc_cb_t out_proc_cb, hg_rpc_cb_t rpc_cb,
    uint32_t mplex_id, ABT_pool pool)
477
{
478 479 480
    struct mplex_key key;
    struct mplex_element *element;
    hg_id_t id;
481

482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506
    id = margo_register_name(mid, func_name, in_proc_cb, out_proc_cb, rpc_cb);
    if(id <= 0)
        return(0);

    /* nothing to do, we'll let the handler pool take this directly */
    if(mplex_id == MARGO_DEFAULT_MPLEX_ID)
        return(id);

    memset(&key, 0, sizeof(key));
    key.id = id;
    key.mplex_id = mplex_id;

    HASH_FIND(hh, mid->mplex_table, &key, sizeof(key), element);
    if(element)
        return(id);

    element = malloc(sizeof(*element));
    if(!element)
        return(0);
    element->key = key;
    element->pool = pool;

    HASH_ADD(hh, mid->mplex_table, key, sizeof(key), element);

    return(id);
507 508
}

509 510
hg_return_t margo_registered_name(margo_instance_id mid, const char *func_name,
    hg_id_t *id, hg_bool_t *flag)
511
{
512
    return(HG_Registered_name(mid->hg_class, func_name, id, flag));
513 514
}

515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 538 539 540 541 542 543 544 545
hg_return_t margo_registered_name_mplex(margo_instance_id mid, const char *func_name,
    uint32_t mplex_id, hg_id_t *id, hg_bool_t *flag)
{
    hg_bool_t b;
    hg_return_t ret = margo_registered_name(mid, func_name, id, &b);
    if(ret != HG_SUCCESS) 
        return ret;
    if((!b) || (!mplex_id)) {
        *flag = b;
        return ret;
    }

    struct mplex_key key;
    struct mplex_element *element;

    memset(&key, 0, sizeof(key));
    key.id = *id;
    key.mplex_id = mplex_id;

    HASH_FIND(hh, mid->mplex_table, &key, sizeof(key), element);
    if(!element) {
        *flag = 0;
        return HG_SUCCESS;
    }

    assert(element->key.id == *id && element->key.mplex_id == mplex_id);

    *flag = 1;
    return HG_SUCCESS;
}

546 547 548 549 550 551 552
hg_return_t margo_register_data(
    margo_instance_id mid,
    hg_id_t id,
    void *data,
    void (*free_callback)(void *)) 
{
	struct margo_rpc_data* margo_data 
553
		= (struct margo_rpc_data*) HG_Registered_data(mid->hg_class, id);
554
	if(!margo_data) return HG_OTHER_ERROR;
555 556 557
    if(margo_data->user_data && margo_data->user_free_callback) {
        (margo_data->user_free_callback)(margo_data->user_data);
    }
558 559 560 561 562 563 564 565 566 567 568 569 570
	margo_data->user_data = data;
	margo_data->user_free_callback = free_callback;
	return HG_SUCCESS;
}

void* margo_registered_data(margo_instance_id mid, hg_id_t id)
{
	struct margo_rpc_data* data
		= (struct margo_rpc_data*) HG_Registered_data(margo_get_class(mid), id);
	if(!data) return NULL;
	else return data->user_data;
}

571 572 573 574
hg_return_t margo_registered_disable_response(
    margo_instance_id mid,
    hg_id_t id,
    int disable_flag)
575
{
576
    return(HG_Registered_disable_response(mid->hg_class, id, disable_flag));
577
}
578

579
struct lookup_cb_evt
580
{
581
    hg_return_t hret;
582 583 584 585 586 587
    hg_addr_t addr;
};

static hg_return_t margo_addr_lookup_cb(const struct hg_cb_info *info)
{
    struct lookup_cb_evt evt;
588
    evt.hret = info->ret;
589
    evt.addr = info->info.lookup.addr;
Matthieu Dorier's avatar
Matthieu Dorier committed
590
    ABT_eventual eventual = (ABT_eventual)(info->arg);
591 592

    /* propagate return code out through eventual */
Matthieu Dorier's avatar
Matthieu Dorier committed
593
    ABT_eventual_set(eventual, &evt, sizeof(evt));
594

595 596 597
    return(HG_SUCCESS);
}

598 599 600 601
hg_return_t margo_addr_lookup(
    margo_instance_id mid,
    const char   *name,
    hg_addr_t    *addr)
602
{
603
    hg_return_t hret;
604 605 606
    struct lookup_cb_evt *evt;
    ABT_eventual eventual;
    int ret;
607

608 609 610 611 612 613
    ret = ABT_eventual_create(sizeof(*evt), &eventual);
    if(ret != 0)
    {
        return(HG_NOMEM_ERROR);        
    }

614
    hret = HG_Addr_lookup(mid->hg_context, margo_addr_lookup_cb,
Matthieu Dorier's avatar
Matthieu Dorier committed
615
        (void*)eventual, name, HG_OP_ID_IGNORE);
616
    if(hret == HG_SUCCESS)
617 618 619
    {
        ABT_eventual_wait(eventual, (void**)&evt);
        *addr = evt->addr;
620
        hret = evt->hret;
621 622 623 624
    }

    ABT_eventual_free(&eventual);

625
    return(hret);
626 627 628 629 630
}

hg_return_t margo_addr_free(
    margo_instance_id mid,
    hg_addr_t addr)
631
{
632 633
    return(HG_Addr_free(mid->hg_class, addr));
}
634

635 636 637 638 639
hg_return_t margo_addr_self(
    margo_instance_id mid,
    hg_addr_t *addr)
{
    return(HG_Addr_self(mid->hg_class, addr));
640 641
}

642 643 644 645 646 647 648 649 650
hg_return_t margo_addr_dup(
    margo_instance_id mid,
    hg_addr_t addr,
    hg_addr_t *new_addr)
{
    return(HG_Addr_dup(mid->hg_class, addr, new_addr));
}

hg_return_t margo_addr_to_string(
651
    margo_instance_id mid,
652 653 654 655 656 657 658 659 660 661
    char *buf,
    hg_size_t *buf_size,
    hg_addr_t addr)
{
    return(HG_Addr_to_string(mid->hg_class, buf, buf_size, addr));
}

hg_return_t margo_create(margo_instance_id mid, hg_addr_t addr,
    hg_id_t id, hg_handle_t *handle)
{
662
    hg_return_t hret = HG_OTHER_ERROR;
Shane Snyder's avatar
Shane Snyder committed
663 664 665 666 667 668 669 670

    /* look for a handle to reuse */
    hret = margo_handle_cache_get(mid, addr, id, handle);
    if(hret != HG_SUCCESS)
    {
        /* else try creating a new handle */
        hret = HG_Create(mid->hg_context, addr, id, handle);
    }
671

Shane Snyder's avatar
Shane Snyder committed
672
    return hret;
673 674
}

675
hg_return_t margo_destroy(hg_handle_t handle)
676
{
677
    margo_instance_id mid;
678
    hg_return_t hret = HG_OTHER_ERROR;
Shane Snyder's avatar
Shane Snyder committed
679

680 681 682
    /* use the handle to get the associated mid */
    mid = margo_hg_handle_get_instance(handle);

Shane Snyder's avatar
Shane Snyder committed
683 684 685 686 687 688 689
    /* recycle this handle if it came from the handle cache */
    hret = margo_handle_cache_put(mid, handle);
    if(hret != HG_SUCCESS)
    {
        /* else destroy the handle manually */
        hret = HG_Destroy(handle);
    }
690

Shane Snyder's avatar
Shane Snyder committed
691
    return hret;
692 693 694 695 696
}

static hg_return_t margo_cb(const struct hg_cb_info *info)
{
    hg_return_t hret = info->ret;
Matthieu Dorier's avatar
Matthieu Dorier committed
697
    ABT_eventual eventual = (ABT_eventual)(info->arg);
698 699

    /* propagate return code out through eventual */
Matthieu Dorier's avatar
Matthieu Dorier committed
700
    ABT_eventual_set(eventual, &hret, sizeof(hret));
701 702 703 704 705 706 707
    
    return(HG_SUCCESS);
}

hg_return_t margo_forward(
    hg_handle_t handle,
    void *in_struct)
Matthieu Dorier's avatar
Matthieu Dorier committed
708 709 710 711 712 713 714 715 716 717 718 719 720
{
	hg_return_t hret;
	margo_request req;
	hret = margo_iforward(handle, in_struct, &req);
	if(hret != HG_SUCCESS) 
		return hret;
	return margo_wait(req);
}

hg_return_t margo_iforward(
    hg_handle_t handle,
    void *in_struct,
    margo_request* req)
721 722
{
    hg_return_t hret = HG_TIMEOUT;
723
    ABT_eventual eventual;
724
    int ret;
725 726 727 728 729 730 731

    ret = ABT_eventual_create(sizeof(hret), &eventual);
    if(ret != 0)
    {
        return(HG_NOMEM_ERROR);        
    }

Matthieu Dorier's avatar
Matthieu Dorier committed
732
    *req = eventual;
733

Matthieu Dorier's avatar
Matthieu Dorier committed
734
    return HG_Forward(handle, margo_cb, (void*)eventual, in_struct);
Matthieu Dorier's avatar
Matthieu Dorier committed
735
}
736

Matthieu Dorier's avatar
Matthieu Dorier committed
737 738 739 740
hg_return_t margo_wait(margo_request req)
{
	hg_return_t* waited_hret;
	hg_return_t  hret;
741

Matthieu Dorier's avatar
Matthieu Dorier committed
742 743 744 745
    ABT_eventual_wait(req, (void**)&waited_hret);
	hret = *waited_hret;
    ABT_eventual_free(&req);
	
746
    return(hret);
747 748
}

Matthieu Dorier's avatar
Matthieu Dorier committed
749 750 751 752 753
int margo_test(margo_request req, int* flag)
{
    return ABT_eventual_test(req, NULL, flag);
}

754 755 756 757
typedef struct
{
    hg_handle_t handle;
} margo_forward_timeout_cb_dat;
758

759 760 761 762 763 764 765 766 767 768 769
static void margo_forward_timeout_cb(void *arg)
{
    margo_forward_timeout_cb_dat *timeout_cb_dat =
        (margo_forward_timeout_cb_dat *)arg;

    /* cancel the Mercury op if the forward timed out */
    HG_Cancel(timeout_cb_dat->handle);
    return;
}

hg_return_t margo_forward_timed(
770
    hg_handle_t handle,
771 772
    void *in_struct,
    double timeout_ms)
773 774
{
    int ret;
775
    hg_return_t hret;
776
    margo_instance_id mid;
777
    ABT_eventual eventual;
778
    hg_return_t* waited_hret;
779 780
    margo_timer_t forward_timer;
    margo_forward_timeout_cb_dat timeout_cb_dat;
781 782 783 784 785 786 787

    ret = ABT_eventual_create(sizeof(hret), &eventual);
    if(ret != 0)
    {
        return(HG_NOMEM_ERROR);        
    }

788 789 790
    /* use the handle to get the associated mid */
    mid = margo_hg_handle_get_instance(handle);

791 792 793 794 795
    /* set a timer object to expire when this forward times out */
    timeout_cb_dat.handle = handle;
    margo_timer_init(mid, &forward_timer, margo_forward_timeout_cb,
        &timeout_cb_dat, timeout_ms);

Matthieu Dorier's avatar
Matthieu Dorier committed
796
    hret = HG_Forward(handle, margo_cb, (void*)eventual, in_struct);
797
    if(hret == HG_SUCCESS)
Jonathan Jenkins's avatar
Jonathan Jenkins committed
798 799 800 801 802
    {
        ABT_eventual_wait(eventual, (void**)&waited_hret);
        hret = *waited_hret;
    }

803 804 805 806 807 808 809 810
    /* convert HG_CANCELED to HG_TIMEOUT to indicate op timed out */
    if(hret == HG_CANCELED)
        hret = HG_TIMEOUT;

    /* remove timer if it is still in place (i.e., not timed out) */
    if(hret != HG_TIMEOUT)
        margo_timer_destroy(mid, &forward_timer);

Jonathan Jenkins's avatar
Jonathan Jenkins committed
811 812 813 814 815 816 817 818
    ABT_eventual_free(&eventual);

    return(hret);
}

hg_return_t margo_respond(
    hg_handle_t handle,
    void *out_struct)
Matthieu Dorier's avatar
Matthieu Dorier committed
819 820 821 822 823 824 825 826 827 828 829 830 831
{
    hg_return_t hret;
    margo_request req;
    hret = margo_irespond(handle,out_struct,&req);
    if(hret != HG_SUCCESS)
        return hret;
    return margo_wait(req);
}

hg_return_t margo_irespond(
    hg_handle_t handle,
    void *out_struct,
    margo_request* req)
Jonathan Jenkins's avatar
Jonathan Jenkins committed
832 833 834 835
{
    ABT_eventual eventual;
    int ret;

Matthieu Dorier's avatar
Matthieu Dorier committed
836
    ret = ABT_eventual_create(sizeof(hg_return_t), &eventual);
Jonathan Jenkins's avatar
Jonathan Jenkins committed
837 838 839 840 841
    if(ret != 0)
    {
        return(HG_NOMEM_ERROR);
    }

Matthieu Dorier's avatar
Matthieu Dorier committed
842
    *req = eventual;
843

Matthieu Dorier's avatar
Matthieu Dorier committed
844
    return HG_Respond(handle, margo_cb, (void*)eventual, out_struct);
845 846
}

847 848 849 850 851 852 853
hg_return_t margo_bulk_create(
    margo_instance_id mid,
    hg_uint32_t count,
    void **buf_ptrs,
    const hg_size_t *buf_sizes,
    hg_uint8_t flags,
    hg_bulk_t *handle)
854
{
855 856 857
    return(HG_Bulk_create(mid->hg_class, count,
        buf_ptrs, buf_sizes, flags, handle));
}
858

859 860 861 862
hg_return_t margo_bulk_free(
    hg_bulk_t handle)
{
    return(HG_Bulk_free(handle));
863 864
}

865 866 867 868 869 870 871 872
hg_return_t margo_bulk_deserialize(
    margo_instance_id mid,
    hg_bulk_t *handle,
    const void *buf,
    hg_size_t buf_size)
{
    return(HG_Bulk_deserialize(mid->hg_class, handle, buf, buf_size));
}
873

874
hg_return_t margo_bulk_transfer(
875
    margo_instance_id mid,
876
    hg_bulk_op_t op,
877
    hg_addr_t origin_addr,
878 879 880 881
    hg_bulk_t origin_handle,
    size_t origin_offset,
    hg_bulk_t local_handle,
    size_t local_offset,
882
    size_t size)
Matthieu Dorier's avatar
Matthieu Dorier committed
883 884 885 886
{  
    margo_request req;
    hg_return_t hret = margo_bulk_itransfer(mid,op,origin_addr,
                          origin_handle, origin_offset, local_handle,
Matthieu Dorier's avatar
Matthieu Dorier committed
887
                          local_offset, size, &req);
Matthieu Dorier's avatar
Matthieu Dorier committed
888 889 890 891 892 893 894 895 896 897 898 899 900 901 902
    if(hret != HG_SUCCESS)
        return hret;
    return margo_wait(req);
}

hg_return_t margo_bulk_itransfer(
    margo_instance_id mid,
    hg_bulk_op_t op,
    hg_addr_t origin_addr,
    hg_bulk_t origin_handle,
    size_t origin_offset,
    hg_bulk_t local_handle,
    size_t local_offset,
    size_t size,
    margo_request* req)
903 904 905 906 907 908 909 910 911 912 913 914
{
    hg_return_t hret = HG_TIMEOUT;
    hg_return_t *waited_hret;
    ABT_eventual eventual;
    int ret;

    ret = ABT_eventual_create(sizeof(hret), &eventual);
    if(ret != 0)
    {
        return(HG_NOMEM_ERROR);        
    }

Matthieu Dorier's avatar
Matthieu Dorier committed
915
    *req = eventual;
916

Matthieu Dorier's avatar
Matthieu Dorier committed
917 918
    hret = HG_Bulk_transfer(mid->hg_context, margo_cb,
        (void*)eventual, op, origin_addr, origin_handle, origin_offset, local_handle,
919
        local_offset, size, HG_OP_ID_IGNORE);
920 921 922 923

    return(hret);
}

924 925 926 927
typedef struct
{
    ABT_mutex mutex;
    ABT_cond cond;
Shane Snyder's avatar
Shane Snyder committed
928
    char is_asleep;
929 930 931 932 933 934 935 936 937
} margo_thread_sleep_cb_dat;

static void margo_thread_sleep_cb(void *arg)
{
    margo_thread_sleep_cb_dat *sleep_cb_dat =
        (margo_thread_sleep_cb_dat *)arg;

    /* wake up the sleeping thread */
    ABT_mutex_lock(sleep_cb_dat->mutex);
938
    sleep_cb_dat->is_asleep = 0;
939 940 941 942 943 944 945
    ABT_cond_signal(sleep_cb_dat->cond);
    ABT_mutex_unlock(sleep_cb_dat->mutex);

    return;
}

void margo_thread_sleep(
946
    margo_instance_id mid,
947 948 949 950 951 952 953 954
    double timeout_ms)
{
    margo_timer_t sleep_timer;
    margo_thread_sleep_cb_dat sleep_cb_dat;

    /* set data needed for sleep callback */
    ABT_mutex_create(&(sleep_cb_dat.mutex));
    ABT_cond_create(&(sleep_cb_dat.cond));
955
    sleep_cb_dat.is_asleep = 1;
956 957

    /* initialize the sleep timer */
958
    margo_timer_init(mid, &sleep_timer, margo_thread_sleep_cb,
959 960 961 962
        &sleep_cb_dat, timeout_ms);

    /* yield thread for specified timeout */
    ABT_mutex_lock(sleep_cb_dat.mutex);
963 964
    while(sleep_cb_dat.is_asleep)
        ABT_cond_wait(sleep_cb_dat.cond, sleep_cb_dat.mutex);
965 966
    ABT_mutex_unlock(sleep_cb_dat.mutex);

Jonathan Jenkins's avatar
Jonathan Jenkins committed
967 968 969 970
    /* clean up */
    ABT_mutex_free(&sleep_cb_dat.mutex);
    ABT_cond_free(&sleep_cb_dat.cond);

971 972 973
    return;
}

974
int margo_get_handler_pool(margo_instance_id mid, ABT_pool* pool)
975
{
976 977 978 979 980 981
    if(mid) {
        *pool = mid->handler_pool;
        return 0;
    } else {
        return -1;
    }
982
}
983

984 985 986 987
hg_context_t* margo_get_context(margo_instance_id mid)
{
    return(mid->hg_context);
}
988

989 990 991
hg_class_t* margo_get_class(margo_instance_id mid)
{
    return(mid->hg_class);
992
}
Philip Carns's avatar
Philip Carns committed
993

994
margo_instance_id margo_hg_handle_get_instance(hg_handle_t h)
995
{
996 997
	const struct hg_info* info = HG_Get_info(h);
	if(!info) return MARGO_INSTANCE_NULL;
998 999 1000 1001 1002
    return margo_hg_info_get_instance(info);
}

margo_instance_id margo_hg_info_get_instance(const struct hg_info *info)
{
1003 1004 1005 1006
	struct margo_rpc_data* data = 
		(struct margo_rpc_data*) HG_Registered_data(info->hg_class, info->id);
	if(!data) return MARGO_INSTANCE_NULL;
	return data->mid;
1007 1008
}

Philip Carns's avatar
Philip Carns committed
1009 1010 1011 1012 1013 1014 1015 1016 1017 1018 1019 1020 1021 1022 1023 1024
int margo_lookup_mplex(margo_instance_id mid, hg_id_t id, uint32_t mplex_id, ABT_pool *pool)
{
    struct mplex_key key;
    struct mplex_element *element;

    if(!mplex_id)
    {
        *pool = mid->handler_pool;
        return(0);
    }

    memset(&key, 0, sizeof(key));
    key.id = id;
    key.mplex_id = mplex_id;

    HASH_FIND(hh, mid->mplex_table, &key, sizeof(key), element);
1025 1026 1027 1028 1029 1030
    if(!element) {
        if(mplex_id == 0) // element does not exist and mplex is 0, return default handler
            *pool = mid->handler_pool;
        else // otherwise it is an error
            return(-1);
    }
Philip Carns's avatar
Philip Carns committed
1031

Philip Carns's avatar
Philip Carns committed
1032 1033
    assert(element->key.id == id && element->key.mplex_id == mplex_id);

Philip Carns's avatar
Philip Carns committed
1034 1035 1036 1037 1038
    *pool = element->pool;

    return(0);
}

1039 1040 1041 1042 1043 1044 1045 1046 1047 1048 1049 1050 1051 1052 1053 1054 1055 1056 1057 1058 1059 1060 1061 1062 1063 1064 1065 1066 1067 1068 1069 1070 1071 1072 1073 1074 1075 1076 1077 1078 1079
int margo_register_data_mplex(margo_instance_id mid, hg_id_t id, uint32_t mplex_id, void* data, void (*free_callback)(void *))
{
    struct mplex_key key;
    struct mplex_element *element;

    memset(&key, 0, sizeof(key));
    key.id = id;
    key.mplex_id = mplex_id;

    HASH_FIND(hh, mid->mplex_table, &key, sizeof(key), element);
    if(!element)
        return -1;

    assert(element->key.id == id && element->key.mplex_id == mplex_id);

    if(element->user_data && element->user_free_callback)
        (element->user_free_callback)(element->user_data);

    element->user_data = data;
    element->user_free_callback = free_callback;

    return(0);
}

void* margo_registered_data_mplex(margo_instance_id mid, hg_id_t id, uint32_t mplex_id)
{
    struct mplex_key key;
    struct mplex_element *element;

    memset(&key, 0, sizeof(key));
    key.id = id;
    key.mplex_id = mplex_id;

    HASH_FIND(hh, mid->mplex_table, &key, sizeof(key), element);
    if(!element)
        return NULL;

    assert(element->key.id == id && element->key.mplex_id == mplex_id);

    return element->user_data;
}
1080
static void margo_rpc_data_free(void* ptr)
Philip Carns's avatar
Philip Carns committed
1081
{
1082 1083 1084 1085 1086 1087
	struct margo_rpc_data* data = (struct margo_rpc_data*) ptr;
	if(data->user_data && data->user_free_callback) {
		data->user_free_callback(data->user_data);
	}
	free(ptr);
}
Philip Carns's avatar
Philip Carns committed
1088

1089 1090 1091