70 for (
size_t i = 0; i < num_entries; ++i)
117 size_t *stripes_array_size,
128 Xt_int start = (
Xt_int)(0 + dist_dir_rank * local_interval);
130 long long start_correction
131 = (
long long)local_index_range_lbound - (
long long)start;
133 = (
Xt_int)((start_correction
135 & (
long long)(global_interval - 1)))
136 / (
long long)global_interval);
137 start = (
Xt_int)(start + corr_steps * global_interval);
143 + (((
long long)local_index_range_ubound - (
long long)start)
144 / global_interval) * global_interval);
146 = (
Xt_int)(start + local_interval > local_index_range_lbound);
147 num_stripes = (int)(((
long long)end - (
long long)start)/global_interval)
148 + (
int)use_start_stripe;
155 struct Xt_stripe *restrict stripes = *stripes_;
158 for (
int j = 0; j < num_stripes; ++j) {
160 stripes[j].start = (
Xt_int)(start + j * global_interval);
161 stripes[j].stride = 1;
162 stripes[j].nstrides = (int)local_interval;
171 void *restrict send_buffer,
172 size_t send_size_asize,
size_t send_size_entry,
173 int tag,
MPI_Comm comm,
int rank_lim,
174 MPI_Request *restrict requests,
175 const int (*send_size)[send_size_asize])
181 for (
size_t rank = 0; rank < (size_t)rank_lim; ++rank)
182 if (send_size[rank][send_size_entry] > 0) {
183 xt_mpi_call(MPI_Isend((
char *)send_buffer + offset,
184 send_size[rank][send_size_entry],
185 MPI_PACKED, (
int)rank, tag,
186 comm, requests + reqOfs),
189 offset += (size_t)send_size[rank][send_size_entry];
191 return (
struct Xt_xmdd_txstat){ .bytes = offset, .num_msg = reqOfs };
196 const struct dist_dir *dst_dist_dir,
197 struct isect **src_dst_intersections)
199 struct isect (*src_dst_intersections_)
200 = (*src_dst_intersections)
203 *
sizeof(**src_dst_intersections));
204 size_t isect_fill = 0;
206 *restrict entries_dst = dst_dist_dir->
entries;
207 size_t num_entries_src = (size_t)src_dist_dir->
num_entries,
208 num_entries_dst = (
size_t)dst_dist_dir->
num_entries;
209 for (
size_t i = 0; i < num_entries_src; ++i)
210 for (
size_t j = 0; j < num_entries_dst; ++j)
215 src_dst_intersections_[isect_fill]
219 .idxlist = intersection };
224 *src_dst_intersections
226 isect_fill *
sizeof (*src_dst_intersections_));
234 size_t num_intersections,
235 const struct isect *restrict src_dst_intersections,
236 bool isect_idxlist_delete,
237 void *buffer,
int buf_size,
int *position,
MPI_Comm comm)
239 int prev_send_rank = -1;
240 int num_send_indices_requests = 0;
241 size_t origin = 1 - target;
242 for (
size_t i = 0; i < num_intersections; ++i)
245 int send_rank = src_dst_intersections[i].rank[target];
246 num_send_indices_requests += send_rank != prev_send_rank;
247 prev_send_rank = send_rank;
251 1, MPI_INT, buffer, buf_size, position,
255 buf_size, position, comm);
257 if (isect_idxlist_delete)
260 return num_send_indices_requests;
269 - (((csx)b)->start > ((csx)a)->start);
283 struct Xt_com_list *restrict entries = (*dist_dir_results)->entries;
284 size_t num_isect_agg = 0;
286 size_t i = 0, num_shards = (size_t)(*dist_dir_results)->num_entries;
287 while (i < num_shards) {
288 int rank = entries[i].rank;
293 while (j < num_shards && entries[j].rank == rank);
295 struct Xt_stripe *restrict stripes = NULL;
296 size_t num_stripes = 0;
298 struct Xt_stripe *stripes_of_intersection;
299 int num_stripes_of_intersection;
301 &stripes_of_intersection,
302 &num_stripes_of_intersection);
306 (num_stripes + (
size_t)num_stripes_of_intersection)
307 *
sizeof (*stripes));
308 memcpy(stripes + num_stripes, stripes_of_intersection,
309 (
size_t)num_stripes_of_intersection *
sizeof (*stripes));
310 free(stripes_of_intersection);
312 stripes = stripes_of_intersection;
313 num_stripes += (size_t)num_stripes_of_intersection;
315 qsort(stripes, num_stripes,
sizeof (*stripes),
stripe_cmp);
318 entries[num_isect_agg].rank =
rank;
321 (*dist_dir_results)->num_entries = (int)num_isect_agg;
323 + (
size_t)num_isect_agg
331 const struct isect *a = a_, *b = b_;
339 const struct isect *a = a_, *b = b_;
349 return a->
rank - b->rank;
357 xt_mpi_call(MPI_Comm_test_inter(comm, &is_inter), comm);
362 MPI_Comm merge_comm, local_intra_comm;
363 MPI_Group local_group;
364 xt_mpi_call(MPI_Comm_group(comm, &local_group), comm);
365 xt_mpi_call(MPI_Intercomm_merge(comm, 0, &merge_comm), comm);
366 xt_mpi_call(MPI_Comm_create(merge_comm, local_group, &local_intra_comm),
373 comm, local_intra_comm);
375 xt_mpi_call(MPI_Comm_free(&local_intra_comm), local_intra_comm);
int xt_idxlist_get_num_indices(Xt_idxlist idxlist)
Uitlity functions for creation of distributed directories.
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_)
Xt_idxlist xt_idxstripes_prealloc_new(const struct Xt_stripe *stripes, int num_stripes)
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
struct Xt_xmdd_txstat xt_xmap_dist_dir_send_intersections(void *restrict send_buffer, size_t send_size_asize, size_t send_size_entry, int tag, MPI_Comm comm, int rank_lim, MPI_Request *restrict requests, const int(*send_size)[send_size_asize])
void xt_mpi_comm_mark_exclusive(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_dist_dir_new(Xt_idxlist src_idxlist, Xt_idxlist dst_idxlist, MPI_Comm comm)
void xt_idxlist_get_index_stripes(Xt_idxlist idxlist, struct Xt_stripe **stripes, int *num_stripes)
#define xrealloc(ptr, size)
Xt_xmap xt_xmap_dist_dir_intercomm_new(Xt_idxlist src_idxlist, Xt_idxlist dst_idxlist, MPI_Comm inter_comm, MPI_Comm intra_comm)
void xt_xmap_dist_dir_same_rank_merge(struct dist_dir **dist_dir_results)
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_)
#define ENSURE_ARRAY_SIZE(arrayp, curr_array_size, req_size)
Xt_int local_index_range_lbound
static int stripe_cmp(const void *a, const void *b)
Xt_xmap xt_xmap_dist_dir_intracomm_new(Xt_idxlist src_idxlist, Xt_idxlist dst_idxlist, MPI_Comm comm)
struct Xt_com_list entries[]
Xt_int local_index_range_ubound
static long long llsign_mask(long long x)
#define xt_mpi_call(call, comm)
Xt_idxlist xt_idxstripes_new(struct Xt_stripe const *stripes, int num_stripes)
void xt_idxlist_pack(Xt_idxlist idxlist, void *buffer, int buffer_size, int *position, MPI_Comm comm)
Xt_idxlist xt_idxlist_get_intersection(Xt_idxlist idxlist_src, Xt_idxlist idxlist_dst)