81 bool *stripify,
Xt_int interval_size[2],
83 int comm_size,
int remote_size,
84 int comm_rank,
int tag_offset_inter)
88 unsigned long long local_vals[2], global_sums[2][2];
91 local_vals[0] = num_indices_src;
95 xt_mpi_call(MPI_Allreduce(local_vals, global_sums[0], 2,
96 MPI_UNSIGNED_LONG_LONG, MPI_SUM, intra_comm),
104 if (comm_rank == 0) {
106 xt_mpi_call(MPI_Sendrecv(global_sums[0], 2, MPI_UNSIGNED_LONG_LONG, 0, tag,
107 global_sums[1], 2, MPI_UNSIGNED_LONG_LONG, 0, tag,
108 inter_comm, MPI_STATUS_IGNORE), inter_comm);
110 xt_mpi_call(MPI_Bcast(global_sums[1], 2, MPI_UNSIGNED_LONG_LONG,
111 0, intra_comm), intra_comm);
112 *stripify = (global_sums[0][1] > 0 || global_sums[1][1] > 0);
114 = (
Xt_int)(((global_sums[0][0] + (
unsigned)comm_size - 1)
115 / (unsigned)comm_size) * (unsigned)comm_size);
117 = (
Xt_int)(((global_sums[1][0] + (
unsigned)remote_size - 1)
118 / (
unsigned)remote_size) * (
unsigned)remote_size);
144 if (local_index_range_lbound <= local_index_range_ubound) {
145 int send_buffer_size = 0;
147 =
xmalloc((
size_t)comm_size *
sizeof(*send_list));
149 size_t stripes_array_size = 0;
152 size_t first_overlapping_bucket = 0;
154 if (local_index_range_lbound >= 0
155 && (local_index_range_ubound < global_interval)) {
156 first_overlapping_bucket
157 = (size_t)(local_index_range_lbound / local_interval);
158 for (
size_t i = 0; i < first_overlapping_bucket; ++i)
162 size_t start_of_non_overlapping_bucket_suffix
163 = (size_t)(((
long long)local_index_range_ubound + local_interval - 1)
164 / local_interval) + 1;
165 if (local_index_range_lbound < 0
166 || start_of_non_overlapping_bucket_suffix > (
size_t)comm_size)
167 start_of_non_overlapping_bucket_suffix = (
size_t)comm_size;
168 size_t i = first_overlapping_bucket;
169 for (; i < (size_t)start_of_non_overlapping_bucket_suffix; ++i) {
172 &stripes, &stripes_array_size, (
int)i);
177 send_list[num_msg] = isect2send;
190 send_buffer_size += send_size[i][
SEND_SIZE];
192 for (; i < (size_t)comm_size; ++i)
196 = *send_buffer_ =
xrealloc(stripes, (
size_t)send_buffer_size);
199 for (i = 0; i < num_msg; ++i) {
207 memset(send_size, 0, (
size_t)comm_size *
sizeof (*send_size));
227 static struct bucket_params
235 return (
struct bucket_params){
247 size_t tx_num = 0, size_sum = 0;
248 for (
size_t i = 0; i < (size_t)comm_size; ++i)
251 size_sum += (size_t)tx_size;
253 if (counts) counts[tx_num] = sizes[i][
SEND_NUM];
263 void **send_buffer,
int send_size[][SEND_SIZE_ASIZE],
267 struct bucket_params bucket_params
271 &bucket_params, idxlist, send_size, send_buffer, comm, comm_size);
272 xt_mpi_call(MPI_Alltoall(send_size, SEND_SIZE_ASIZE, MPI_INT,
273 recv_size, SEND_SIZE_ASIZE, MPI_INT, comm), comm);
276 assert(num_msg == tx_stat[1].num_msg);
279 typedef int (*
tx_fp)(
void *, int, MPI_Datatype, int, int,
284 unsigned char *buffer, MPI_Request *requests,
285 int tag, MPI_Comm comm,
tx_fp tx_op)
288 for (
size_t i = 0; i < num_msg; ++i)
292 count, MPI_PACKED, rank, tag, comm, requests + i), comm);
293 ofs += (size_t)count;
300 void *recv_buffer, MPI_Request *requests,
301 int tag, MPI_Comm comm)
310 void *send_buffer, MPI_Request *requests,
311 int tag, MPI_Comm comm)
325 size_t num_msg = tx_stat.
num_msg, buf_size = tx_stat.
bytes;
326 struct dist_dir *restrict dist_dir_
327 = *dist_dir =
xmalloc(
sizeof (*dist_dir_)
328 +
sizeof (*dist_dir_->entries) * num_msg);
329 dist_dir_->num_entries = (int)num_msg;
331 for (
size_t i = 0; i < num_msg; ++i)
334 dist_dir_->entries[i].rank = rank;
335 dist_dir_->entries[i].list
346 int tag_offset_inter,
int tag_offset_intra,
347 MPI_Comm inter_comm, MPI_Comm intra_comm,
348 int remote_size,
int comm_size,
355 stripify, interval_size,
356 intra_comm, inter_comm,
357 comm_size, remote_size,
358 comm_rank, tag_offset_inter);
359 void *send_buffer_local, *send_buffer_remote;
361 =
xmalloc(((
size_t)comm_size + (
size_t)remote_size)
362 * 2 *
sizeof(*send_size_local)),
368 send_size_local, src_idxlist, interval_size[0],
369 intra_comm, comm_size);
371 send_size_remote, dst_idxlist, interval_size[1],
372 inter_comm, remote_size);
374 size_t num_req = tx_stat_local[0].
num_msg + tx_stat_remote[0].
num_msg 376 MPI_Request *dir_init_requests
377 =
xmalloc(num_req *
sizeof(*dir_init_requests)
378 + tx_stat_local[0].
bytes + tx_stat_remote[0].
bytes);
379 void *recv_buffer_local = dir_init_requests + num_req,
380 *recv_buffer_remote = ((
unsigned char *)recv_buffer_local
381 + tx_stat_local[0].bytes);
383 size_t req_ofs = tx_stat_local[0].
num_msg;
386 recv_buffer_local, dir_init_requests,
387 tag_intra, intra_comm);
391 recv_buffer_remote, dir_init_requests + req_ofs,
392 tag_inter, inter_comm);
393 req_ofs += tx_stat_remote[0].
num_msg;
396 send_buffer_local, dir_init_requests + req_ofs,
397 tag_intra, intra_comm);
398 req_ofs += tx_stat_local[1].
num_msg;
401 send_buffer_remote, dir_init_requests + req_ofs,
402 tag_inter, inter_comm);
404 xt_mpi_call(MPI_Waitall((
int)num_req, dir_init_requests,
405 MPI_STATUSES_IGNORE), inter_comm);
406 free(send_buffer_local);
407 free(send_buffer_remote);
410 recv_buffer_local, src_dist_dir, intra_comm);
413 recv_buffer_remote, dst_dist_dir, inter_comm);
414 free(send_size_local);
415 free(dir_init_requests);
421 const struct isect *restrict src_dst_intersections,
423 MPI_Comm comm,
int comm_size,
426 size_t total_send_size = 0;
427 for (
int i = 0; i < comm_size; ++i)
431 xt_mpi_call(MPI_Pack_size(1, MPI_INT, comm, &rank_pack_size), comm);
433 for (
size_t i = 0; i < num_intersections; ++i)
435 int msg_size = rank_pack_size
437 size_t target_rank = (size_t)src_dst_intersections[i].rank[target];
438 send_size_target[target_rank][
SEND_SIZE] += msg_size;
439 ++(send_size_target[target_rank][
SEND_NUM]);
440 total_send_size += (size_t)msg_size;
442 assert(total_send_size <= INT_MAX);
443 return (
int)total_send_size;
449 struct isect *restrict src_dst_intersections,
452 bool isect_idxlist_delete, MPI_Comm comm,
int comm_size) {
456 src_dst_intersections,
458 comm, comm_size, send_size);
460 void *send_buffer = (*send_buffer_) =
xmalloc((
size_t)total_send_size);
462 qsort(src_dst_intersections, num_intersections,
463 sizeof (src_dst_intersections[0]),
468 target, num_intersections, src_dst_intersections,
469 isect_idxlist_delete, send_buffer, total_send_size, &position, comm);
476 void *restrict recv_buffer,
477 int *restrict entry_counts,
480 size_t num_msg = tx_stat.
num_msg;
481 int buf_size = (int)tx_stat.
bytes;
483 size_t num_entries_sent = 0;
484 for (
size_t i = 0; i <
num_msg; ++i)
485 num_entries_sent += (
size_t)entry_counts[i];
486 *dist_dir =
xmalloc(
sizeof (
struct dist_dir)
487 + (
sizeof (
struct Xt_com_list) * num_entries_sent));
488 (*dist_dir)->num_entries = (int)num_entries_sent;
489 struct Xt_com_list *restrict entries = (*dist_dir)->entries;
490 size_t num_entries = 0;
491 for (
size_t i = 0; i < num_msg; ++i) {
492 size_t num_entries_from_rank = (size_t)entry_counts[i];
493 for (
size_t j = 0; j < num_entries_from_rank; ++j) {
494 xt_mpi_call(MPI_Unpack(recv_buffer, buf_size, &position,
495 &entries[num_entries].rank,
496 1, MPI_INT, comm), comm);
497 entries[num_entries].list =
502 assert(num_entries == num_entries_sent);
510 struct dist_dir **dst_intersections,
514 int tag_offset_inter,
int tag_offset_intra,
515 MPI_Comm inter_comm, MPI_Comm intra_comm) {
517 int comm_size, remote_size, comm_rank;
518 xt_mpi_call(MPI_Comm_size(inter_comm, &comm_size), inter_comm);
519 xt_mpi_call(MPI_Comm_rank(inter_comm, &comm_rank), inter_comm);
520 xt_mpi_call(MPI_Comm_remote_size(inter_comm, &remote_size), inter_comm);
522 struct dist_dir *src_dist_dir, *dst_dist_dir;
525 src_idxlist, dst_idxlist,
526 tag_offset_inter, tag_offset_intra,
527 inter_comm, intra_comm,
528 remote_size, comm_size,
533 =
xmalloc(((
size_t)comm_size + (
size_t)remote_size)
534 * 2U *
sizeof(*send_size_local)),
541 struct isect *src_dst_intersections;
542 size_t num_intersections
544 &src_dst_intersections);
548 void *send_buffer_local, *send_buffer_remote;
549 int num_send_requests_local
551 send_size_local, &send_buffer_local,
553 num_send_requests_remote
555 send_size_remote, &send_buffer_remote,
557 free(src_dst_intersections);
562 intra_comm), intra_comm);
565 inter_comm), inter_comm);
568 int *isect_counts_recv_local
569 =
xmalloc(((
size_t)comm_size + (
size_t)remote_size) *
sizeof (
int)),
570 *isect_counts_recv_remote = isect_counts_recv_local + comm_size;
571 compress_sizes(send_size_local, comm_size, tx_stat_local+1, NULL);
573 isect_counts_recv_local);
574 compress_sizes(send_size_remote, remote_size, tx_stat_remote+1, NULL);
576 isect_counts_recv_remote);
577 assert(tx_stat_local[1].num_msg == (
size_t)num_send_requests_local
578 && tx_stat_remote[1].num_msg == (
size_t)num_send_requests_remote);
580 = (size_t)num_send_requests_local + (
size_t)num_send_requests_remote
582 MPI_Request *requests
583 =
xmalloc(num_requests *
sizeof(*requests)
584 + tx_stat_local[0].
bytes + tx_stat_remote[0].
bytes);
585 void *recv_buf_local = requests + num_requests,
586 *recv_buf_remote = (
unsigned char *)recv_buf_local + tx_stat_local[0].bytes;
587 size_t req_ofs = tx_stat_local[0].
num_msg;
591 recv_buf_local, requests, tag_intra, intra_comm);
595 recv_buf_remote, requests+req_ofs, tag_inter, inter_comm);
596 req_ofs += tx_stat_remote[0].
num_msg;
599 send_buffer_local, requests+req_ofs, tag_intra,
601 req_ofs += tx_stat_local[1].
num_msg;
604 send_buffer_remote, requests+req_ofs, tag_inter,
606 xt_mpi_call(MPI_Waitall((
int)num_requests, requests, MPI_STATUSES_IGNORE),
608 free(send_buffer_local);
609 free(send_buffer_remote);
610 free(send_size_local);
613 isect_counts_recv_local, intra_comm);
615 isect_counts_recv_remote, inter_comm);
617 free(isect_counts_recv_local);
624 MPI_Comm inter_comm_, MPI_Comm intra_comm_)
626 INSTR_DEF(this_instr,
"xt_xmap_dist_dir_intercomm_new")
628 int tag_offset_inter, tag_offset_intra;
632 struct dist_dir *src_intersections, *dst_intersections;
636 src_idxlist, dst_idxlist,
637 tag_offset_inter, tag_offset_intra,
638 inter_comm, intra_comm);
640 Xt_xmap (*xmap_new)(
int num_src_intersections,
642 int num_dst_intersections,
651 src_idxlist, dst_idxlist, inter_comm);
Xt_xmap xt_xmap_dist_dir_intercomm_new(Xt_idxlist src_idxlist, Xt_idxlist dst_idxlist, MPI_Comm inter_comm_, MPI_Comm intra_comm_)
static void compress_sizes(int(*restrict sizes)[SEND_SIZE_ASIZE], int comm_size, struct Xt_xmdd_txstat *tx_stat, int *counts)
Xt_idxlist xt_idxlist_unpack(void *buffer, int buffer_size, int *position, MPI_Comm comm)
int xt_idxlist_get_num_indices(Xt_idxlist idxlist)
static void irecv_intersections(size_t num_msg, const int(*recv_size)[SEND_SIZE_ASIZE], void *recv_buffer, MPI_Request *requests, int tag, MPI_Comm comm)
struct Xt_xmap_ * Xt_xmap
Uitlity functions for creation of distributed directories.
static int send_size_from_intersections(size_t num_intersections, const struct isect *restrict src_dst_intersections, enum xt_xmdd_direction target, MPI_Comm comm, int comm_size, int(*restrict send_size_target)[SEND_SIZE_ASIZE])
int xt_com_list_rank_cmp(const void *a_, const void *b_)
int xt_xmdd_cmp_isect_dst_rank(const void *a_, const void *b_)
size_t xt_xmap_dist_dir_match_src_dst(const struct dist_dir *src_dist_dir, const struct dist_dir *dst_dist_dir, struct isect **src_dst_intersections)
add versions of standard API functions not returning on error
static void isend_intersections(size_t num_msg, const int(*send_size)[SEND_SIZE_ASIZE], void *send_buffer, MPI_Request *requests, int tag, MPI_Comm comm)
void xt_idxlist_delete(Xt_idxlist idxlist)
Xt_idxlist xt_xmap_dist_dir_get_bucket(const struct bucket_params *bucket_params, struct Xt_stripe **stripes_, size_t *stripes_array_size, int dist_dir_rank)
generates the buckets of the distributed directory
Xt_xmap xt_xmap_intersection_new(int num_src_intersections, const struct Xt_com_list src_com[num_src_intersections], int num_dst_intersections, const struct Xt_com_list dst_com[num_dst_intersections], Xt_idxlist src_idxlist, Xt_idxlist dst_idxlist, MPI_Comm comm)
Xt_int xt_idxlist_get_min_index(Xt_idxlist idxlist)
#define xrealloc(ptr, size)
static void generate_distributed_directories(struct dist_dir **src_dist_dir, struct dist_dir **dst_dist_dir, bool *stripify, Xt_idxlist src_idxlist, Xt_idxlist dst_idxlist, int tag_offset_inter, int tag_offset_intra, MPI_Comm inter_comm, MPI_Comm intra_comm, int remote_size, int comm_size, int comm_rank)
exchange map declarations
void xt_xmap_dist_dir_same_rank_merge(struct dist_dir **dist_dir_results)
static Xt_int get_min_idxlist_index(Xt_idxlist l)
Provide non-public declarations common to all index lists.
int xt_xmap_dist_dir_pack_intersections(enum xt_xmdd_direction target, size_t num_intersections, const struct isect *restrict src_dst_intersections, bool isect_idxlist_delete, void *buffer, int buf_size, int *position, MPI_Comm comm)
void xt_xmdd_free_dist_dir(struct dist_dir *dist_dir)
int xt_xmdd_cmp_isect_src_rank(const void *a_, const void *b_)
Xt_int xt_idxlist_get_max_index(Xt_idxlist idxlist)
static void exchange_idxlists(struct dist_dir **src_intersections, struct dist_dir **dst_intersections, bool *stripify, Xt_idxlist src_idxlist, Xt_idxlist dst_idxlist, int tag_offset_inter, int tag_offset_intra, MPI_Comm inter_comm, MPI_Comm intra_comm)
static void create_intersections(struct Xt_xmdd_txstat tx_stat[2], int recv_size[][SEND_SIZE_ASIZE], void **send_buffer, int send_size[][SEND_SIZE_ASIZE], Xt_idxlist idxlist, Xt_int interval_size, MPI_Comm comm, int comm_size)
static int pack_dist_dirs(size_t num_intersections, struct isect *restrict src_dst_intersections, int(*send_size)[SEND_SIZE_ASIZE], void **send_buffer_, enum xt_xmdd_direction target, bool isect_idxlist_delete, MPI_Comm comm, int comm_size)
Xt_int local_index_range_lbound
static void unpack_dist_dir(struct Xt_xmdd_txstat tx_stat, const int(*sizes)[SEND_SIZE_ASIZE], const void *buffer, struct dist_dir **dist_dir, MPI_Comm comm)
Xt_xmap xt_xmap_intersection_ext_new(int num_src_intersections, const struct Xt_com_list src_com[num_src_intersections], int num_dst_intersections, const struct Xt_com_list dst_com[num_dst_intersections], Xt_idxlist src_idxlist, Xt_idxlist dst_idxlist, MPI_Comm comm)
struct Xt_com_list entries[]
Xt_int local_index_range_ubound
static void get_dist_dir_global_interval_size(Xt_idxlist src, Xt_idxlist dst, bool *stripify, Xt_int interval_size[2], MPI_Comm intra_comm, MPI_Comm inter_comm, int comm_size, int remote_size, int comm_rank, int tag_offset_inter)
static Xt_int get_max_idxlist_index(Xt_idxlist l)
MPI_Comm xt_mpi_comm_smart_dup(MPI_Comm comm, int *tag_offset)
static void tx_intersections(size_t num_msg, const int(*sizes)[SEND_SIZE_ASIZE], unsigned char *buffer, MPI_Request *requests, int tag, MPI_Comm comm, tx_fp tx_op)
void xt_mpi_comm_smart_dedup(MPI_Comm *comm, int tag_offset)
static struct bucket_params get_bucket_params(Xt_idxlist idxlist, Xt_int global_interval, int comm_size)
size_t xt_idxlist_get_pack_size(Xt_idxlist idxlist, MPI_Comm comm)
#define xt_mpi_call(call, comm)
static size_t compute_and_pack_bucket_intersections(struct bucket_params *bucket_params, Xt_idxlist idxlist, int(*send_size)[SEND_SIZE_ASIZE], void **send_buffer_, MPI_Comm comm, int comm_size)
int(* tx_fp)(void *, int, MPI_Datatype, int, int, MPI_Comm, MPI_Request *)
void xt_idxlist_pack(Xt_idxlist idxlist, void *buffer, int buffer_size, int *position, MPI_Comm comm)
static void rank_no_send(size_t rank, int(*restrict send_size)[SEND_SIZE_ASIZE])
static void unpack_dist_dir_results(struct Xt_xmdd_txstat tx_stat, struct dist_dir **dist_dir, void *restrict recv_buffer, int *restrict entry_counts, MPI_Comm comm)
Xt_idxlist xt_idxlist_get_intersection(Xt_idxlist idxlist_src, Xt_idxlist idxlist_dst)