starpu_mpi.c 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453
  1. /* StarPU --- Runtime system for heterogeneous multicore architectures.
  2. *
  3. * Copyright (C) 2012,2013,2016,2017 Inria
  4. * Copyright (C) 2010-2018 CNRS
  5. * Copyright (C) 2009-2018 Université de Bordeaux
  6. *
  7. * StarPU is free software; you can redistribute it and/or modify
  8. * it under the terms of the GNU Lesser General Public License as published by
  9. * the Free Software Foundation; either version 2.1 of the License, or (at
  10. * your option) any later version.
  11. *
  12. * StarPU is distributed in the hope that it will be useful, but
  13. * WITHOUT ANY WARRANTY; without even the implied warranty of
  14. * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.
  15. *
  16. * See the GNU Lesser General Public License in COPYING.LGPL for more details.
  17. */
  18. #include <stdlib.h>
  19. #include <limits.h>
  20. #include <starpu_mpi.h>
  21. #include <starpu_mpi_datatype.h>
  22. #include <starpu_mpi_private.h>
  23. #include <starpu_mpi_cache.h>
  24. #include <starpu_profiling.h>
  25. #include <starpu_mpi_stats.h>
  26. #include <starpu_mpi_cache.h>
  27. #include <starpu_mpi_select_node.h>
  28. #include <starpu_mpi_init.h>
  29. #include <common/config.h>
  30. #include <common/thread.h>
  31. #include <datawizard/interfaces/data_interface.h>
  32. #include <datawizard/coherency.h>
  33. #include <core/simgrid.h>
  34. #include <core/task.h>
  35. #include <core/topology.h>
  36. #include <core/workers.h>
  37. #if defined(STARPU_USE_MPI_MPI)
  38. #include <mpi/starpu_mpi_comm.h>
  39. #include <mpi/starpu_mpi_tag.h>
  40. #endif
  41. static void _starpu_mpi_isend_irecv_common(struct _starpu_mpi_req *req, enum starpu_data_access_mode mode, int sequential_consistency)
  42. {
  43. /* Asynchronously request StarPU to fetch the data in main memory: when
  44. * it is available in main memory, _starpu_mpi_submit_ready_request(req) is called and
  45. * the request is actually submitted */
  46. starpu_data_acquire_on_node_cb_sequential_consistency_sync_jobids(req->data_handle, STARPU_MAIN_RAM, mode, _starpu_mpi_submit_ready_request, (void *)req, sequential_consistency, 1, &req->pre_sync_jobid, &req->post_sync_jobid);
  47. }
  48. static struct _starpu_mpi_req *_starpu_mpi_isend_common(starpu_data_handle_t data_handle,
  49. int dest, starpu_mpi_tag_t data_tag, MPI_Comm comm,
  50. unsigned detached, unsigned sync, int prio, void (*callback)(void *), void *arg,
  51. int sequential_consistency)
  52. {
  53. if (_starpu_mpi_fake_world_size != -1)
  54. {
  55. /* Don't actually do the communication */
  56. return NULL;
  57. }
  58. #ifdef STARPU_MPI_PEDANTIC_ISEND
  59. enum starpu_data_access_mode mode = STARPU_RW;
  60. #else
  61. enum starpu_data_access_mode mode = STARPU_R;
  62. #endif
  63. struct _starpu_mpi_req *req = _starpu_mpi_request_fill(
  64. data_handle, dest, data_tag, comm, detached, sync, prio, callback, arg, SEND_REQ, _starpu_mpi_isend_size_func,
  65. sequential_consistency, 0, 0);
  66. _starpu_mpi_req_willpost(req);
  67. if (_starpu_mpi_use_coop_sends && detached == 1 && sync == 0 && callback == NULL)
  68. {
  69. /* It's a send & forget send, we can perhaps optimize its distribution over several nodes */
  70. _starpu_mpi_coop_send(data_handle, req, mode, sequential_consistency);
  71. return req;
  72. }
  73. /* Post normally */
  74. _starpu_mpi_isend_irecv_common(req, mode, sequential_consistency);
  75. return req;
  76. }
  77. int starpu_mpi_isend_prio(starpu_data_handle_t data_handle, starpu_mpi_req *public_req, int dest, starpu_mpi_tag_t data_tag, int prio, MPI_Comm comm)
  78. {
  79. _STARPU_MPI_LOG_IN();
  80. STARPU_MPI_ASSERT_MSG(public_req, "starpu_mpi_isend needs a valid starpu_mpi_req");
  81. struct _starpu_mpi_req *req;
  82. _STARPU_MPI_TRACE_ISEND_COMPLETE_BEGIN(dest, data_tag, 0);
  83. req = _starpu_mpi_isend_common(data_handle, dest, data_tag, comm, 0, 0, prio, NULL, NULL, 1);
  84. _STARPU_MPI_TRACE_ISEND_COMPLETE_END(dest, data_tag, 0);
  85. STARPU_MPI_ASSERT_MSG(req, "Invalid return for _starpu_mpi_isend_common");
  86. *public_req = req;
  87. _STARPU_MPI_LOG_OUT();
  88. return 0;
  89. }
  90. int starpu_mpi_isend(starpu_data_handle_t data_handle, starpu_mpi_req *public_req, int dest, starpu_mpi_tag_t data_tag, MPI_Comm comm)
  91. {
  92. return starpu_mpi_isend_prio(data_handle, public_req, dest, data_tag, 0, comm);
  93. }
  94. int starpu_mpi_isend_detached_prio(starpu_data_handle_t data_handle, int dest, starpu_mpi_tag_t data_tag, int prio, MPI_Comm comm, void (*callback)(void *), void *arg)
  95. {
  96. _STARPU_MPI_LOG_IN();
  97. _starpu_mpi_isend_common(data_handle, dest, data_tag, comm, 1, 0, prio, callback, arg, 1);
  98. _STARPU_MPI_LOG_OUT();
  99. return 0;
  100. }
  101. int starpu_mpi_isend_detached(starpu_data_handle_t data_handle, int dest, starpu_mpi_tag_t data_tag, MPI_Comm comm, void (*callback)(void *), void *arg)
  102. {
  103. return starpu_mpi_isend_detached_prio(data_handle, dest, data_tag, 0, comm, callback, arg);
  104. }
  105. int starpu_mpi_send_prio(starpu_data_handle_t data_handle, int dest, starpu_mpi_tag_t data_tag, int prio, MPI_Comm comm)
  106. {
  107. starpu_mpi_req req;
  108. MPI_Status status;
  109. _STARPU_MPI_LOG_IN();
  110. memset(&status, 0, sizeof(MPI_Status));
  111. starpu_mpi_isend_prio(data_handle, &req, dest, data_tag, prio, comm);
  112. starpu_mpi_wait(&req, &status);
  113. _STARPU_MPI_LOG_OUT();
  114. return 0;
  115. }
  116. int starpu_mpi_send(starpu_data_handle_t data_handle, int dest, starpu_mpi_tag_t data_tag, MPI_Comm comm)
  117. {
  118. return starpu_mpi_send_prio(data_handle, dest, data_tag, 0, comm);
  119. }
  120. int starpu_mpi_issend_prio(starpu_data_handle_t data_handle, starpu_mpi_req *public_req, int dest, starpu_mpi_tag_t data_tag, int prio, MPI_Comm comm)
  121. {
  122. _STARPU_MPI_LOG_IN();
  123. STARPU_MPI_ASSERT_MSG(public_req, "starpu_mpi_issend needs a valid starpu_mpi_req");
  124. struct _starpu_mpi_req *req;
  125. req = _starpu_mpi_isend_common(data_handle, dest, data_tag, comm, 0, 1, prio, NULL, NULL, 1);
  126. STARPU_MPI_ASSERT_MSG(req, "Invalid return for _starpu_mpi_isend_common");
  127. *public_req = req;
  128. _STARPU_MPI_LOG_OUT();
  129. return 0;
  130. }
  131. int starpu_mpi_issend(starpu_data_handle_t data_handle, starpu_mpi_req *public_req, int dest, starpu_mpi_tag_t data_tag, MPI_Comm comm)
  132. {
  133. return starpu_mpi_issend_prio(data_handle, public_req, dest, data_tag, 0, comm);
  134. }
  135. int starpu_mpi_issend_detached_prio(starpu_data_handle_t data_handle, int dest, starpu_mpi_tag_t data_tag, int prio, MPI_Comm comm, void (*callback)(void *), void *arg)
  136. {
  137. _STARPU_MPI_LOG_IN();
  138. _starpu_mpi_isend_common(data_handle, dest, data_tag, comm, 1, 1, prio, callback, arg, 1);
  139. _STARPU_MPI_LOG_OUT();
  140. return 0;
  141. }
  142. int starpu_mpi_issend_detached(starpu_data_handle_t data_handle, int dest, starpu_mpi_tag_t data_tag, MPI_Comm comm, void (*callback)(void *), void *arg)
  143. {
  144. return starpu_mpi_issend_detached_prio(data_handle, dest, data_tag, 0, comm, callback, arg);
  145. }
  146. struct _starpu_mpi_req *_starpu_mpi_irecv_common(starpu_data_handle_t data_handle, int source, starpu_mpi_tag_t data_tag, MPI_Comm comm, unsigned detached, unsigned sync, void (*callback)(void *), void *arg, int sequential_consistency, int is_internal_req, starpu_ssize_t count)
  147. {
  148. if (_starpu_mpi_fake_world_size != -1)
  149. {
  150. /* Don't actually do the communication */
  151. return NULL;
  152. }
  153. struct _starpu_mpi_req *req = _starpu_mpi_request_fill(data_handle, source, data_tag, comm, detached, sync, 0, callback, arg, RECV_REQ, _starpu_mpi_irecv_size_func, sequential_consistency, is_internal_req, count);
  154. _starpu_mpi_req_willpost(req);
  155. _starpu_mpi_isend_irecv_common(req, STARPU_W, sequential_consistency);
  156. return req;
  157. }
  158. int starpu_mpi_irecv(starpu_data_handle_t data_handle, starpu_mpi_req *public_req, int source, starpu_mpi_tag_t data_tag, MPI_Comm comm)
  159. {
  160. _STARPU_MPI_LOG_IN();
  161. STARPU_MPI_ASSERT_MSG(public_req, "starpu_mpi_irecv needs a valid starpu_mpi_req");
  162. struct _starpu_mpi_req *req;
  163. _STARPU_MPI_TRACE_IRECV_COMPLETE_BEGIN(source, data_tag);
  164. req = _starpu_mpi_irecv_common(data_handle, source, data_tag, comm, 0, 0, NULL, NULL, 1, 0, 0);
  165. _STARPU_MPI_TRACE_IRECV_COMPLETE_END(source, data_tag);
  166. STARPU_MPI_ASSERT_MSG(req, "Invalid return for _starpu_mpi_irecv_common");
  167. *public_req = req;
  168. _STARPU_MPI_LOG_OUT();
  169. return 0;
  170. }
  171. int starpu_mpi_irecv_detached(starpu_data_handle_t data_handle, int source, starpu_mpi_tag_t data_tag, MPI_Comm comm, void (*callback)(void *), void *arg)
  172. {
  173. _STARPU_MPI_LOG_IN();
  174. _starpu_mpi_irecv_common(data_handle, source, data_tag, comm, 1, 0, callback, arg, 1, 0, 0);
  175. _STARPU_MPI_LOG_OUT();
  176. return 0;
  177. }
  178. int starpu_mpi_irecv_detached_sequential_consistency(starpu_data_handle_t data_handle, int source, starpu_mpi_tag_t data_tag, MPI_Comm comm, void (*callback)(void *), void *arg, int sequential_consistency)
  179. {
  180. _STARPU_MPI_LOG_IN();
  181. _starpu_mpi_irecv_common(data_handle, source, data_tag, comm, 1, 0, callback, arg, sequential_consistency, 0, 0);
  182. _STARPU_MPI_LOG_OUT();
  183. return 0;
  184. }
  185. int starpu_mpi_recv(starpu_data_handle_t data_handle, int source, starpu_mpi_tag_t data_tag, MPI_Comm comm, MPI_Status *status)
  186. {
  187. starpu_mpi_req req;
  188. _STARPU_MPI_LOG_IN();
  189. starpu_mpi_irecv(data_handle, &req, source, data_tag, comm);
  190. starpu_mpi_wait(&req, status);
  191. _STARPU_MPI_LOG_OUT();
  192. return 0;
  193. }
  194. int starpu_mpi_wait(starpu_mpi_req *public_req, MPI_Status *status)
  195. {
  196. return _starpu_mpi_wait(public_req, status);
  197. }
  198. int starpu_mpi_test(starpu_mpi_req *public_req, int *flag, MPI_Status *status)
  199. {
  200. return _starpu_mpi_test(public_req, flag, status);
  201. }
  202. int starpu_mpi_barrier(MPI_Comm comm)
  203. {
  204. return _starpu_mpi_barrier(comm);
  205. }
  206. void _starpu_mpi_data_clear(starpu_data_handle_t data_handle)
  207. {
  208. #if defined(STARPU_USE_MPI_MPI)
  209. _starpu_mpi_tag_data_release(data_handle);
  210. #endif
  211. _starpu_mpi_cache_data_clear(data_handle);
  212. free(data_handle->mpi_data);
  213. data_handle->mpi_data = NULL;
  214. }
  215. struct _starpu_mpi_data *_starpu_mpi_data_get(starpu_data_handle_t data_handle)
  216. {
  217. struct _starpu_mpi_data *mpi_data = data_handle->mpi_data;
  218. if (mpi_data)
  219. {
  220. STARPU_ASSERT(mpi_data->magic == 42);
  221. }
  222. else
  223. {
  224. _STARPU_CALLOC(mpi_data, 1, sizeof(struct _starpu_mpi_data));
  225. mpi_data->magic = 42;
  226. mpi_data->node_tag.data_tag = -1;
  227. mpi_data->node_tag.rank = -1;
  228. mpi_data->node_tag.comm = MPI_COMM_WORLD;
  229. _starpu_spin_init(&mpi_data->coop_lock);
  230. data_handle->mpi_data = mpi_data;
  231. _starpu_mpi_cache_data_init(data_handle);
  232. _starpu_data_set_unregister_hook(data_handle, _starpu_mpi_data_clear);
  233. }
  234. return mpi_data;
  235. }
  236. void starpu_mpi_data_register_comm(starpu_data_handle_t data_handle, starpu_mpi_tag_t data_tag, int rank, MPI_Comm comm)
  237. {
  238. struct _starpu_mpi_data *mpi_data = _starpu_mpi_data_get(data_handle);
  239. if (data_tag != -1)
  240. {
  241. #if defined(STARPU_USE_MPI_MPI)
  242. _starpu_mpi_tag_data_register(data_handle, data_tag);
  243. #endif
  244. mpi_data->node_tag.data_tag = data_tag;
  245. }
  246. if (rank != -1)
  247. {
  248. _STARPU_MPI_TRACE_DATA_SET_RANK(data_handle, rank);
  249. mpi_data->node_tag.rank = rank;
  250. mpi_data->node_tag.comm = comm;
  251. #if defined(STARPU_USE_MPI_MPI)
  252. _starpu_mpi_comm_register(comm);
  253. #endif
  254. }
  255. }
  256. void starpu_mpi_data_set_rank_comm(starpu_data_handle_t handle, int rank, MPI_Comm comm)
  257. {
  258. starpu_mpi_data_register_comm(handle, -1, rank, comm);
  259. }
  260. void starpu_mpi_data_set_tag(starpu_data_handle_t handle, starpu_mpi_tag_t data_tag)
  261. {
  262. starpu_mpi_data_register_comm(handle, data_tag, -1, MPI_COMM_WORLD);
  263. }
  264. int starpu_mpi_data_get_rank(starpu_data_handle_t data)
  265. {
  266. STARPU_ASSERT_MSG(data->mpi_data, "starpu_mpi_data_register MUST be called for data %p\n", data);
  267. return ((struct _starpu_mpi_data *)(data->mpi_data))->node_tag.rank;
  268. }
  269. starpu_mpi_tag_t starpu_mpi_data_get_tag(starpu_data_handle_t data)
  270. {
  271. STARPU_ASSERT_MSG(data->mpi_data, "starpu_mpi_data_register MUST be called for data %p\n", data);
  272. return ((struct _starpu_mpi_data *)(data->mpi_data))->node_tag.data_tag;
  273. }
  274. void starpu_mpi_get_data_on_node_detached(MPI_Comm comm, starpu_data_handle_t data_handle, int node, void (*callback)(void*), void *arg)
  275. {
  276. int me, rank, tag;
  277. rank = starpu_mpi_data_get_rank(data_handle);
  278. if (rank == -1)
  279. {
  280. _STARPU_ERROR("StarPU needs to be told the MPI rank of this data, using starpu_mpi_data_register() or starpu_mpi_data_register()\n");
  281. }
  282. starpu_mpi_comm_rank(comm, &me);
  283. if (node == rank)
  284. return;
  285. tag = starpu_mpi_data_get_tag(data_handle);
  286. if (tag == -1)
  287. {
  288. _STARPU_ERROR("StarPU needs to be told the MPI tag of this data, using starpu_mpi_data_register() or starpu_mpi_data_register()\n");
  289. }
  290. if (me == node)
  291. {
  292. _STARPU_MPI_DEBUG(1, "Migrating data %p from %d to %d\n", data_handle, rank, node);
  293. int already_received = _starpu_mpi_cache_received_data_set(data_handle);
  294. if (already_received == 0)
  295. {
  296. _STARPU_MPI_DEBUG(1, "Receiving data %p from %d\n", data_handle, rank);
  297. starpu_mpi_irecv_detached(data_handle, rank, tag, comm, callback, arg);
  298. }
  299. }
  300. else if (me == rank)
  301. {
  302. _STARPU_MPI_DEBUG(1, "Migrating data %p from %d to %d\n", data_handle, rank, node);
  303. int already_sent = _starpu_mpi_cache_sent_data_set(data_handle, node);
  304. if (already_sent == 0)
  305. {
  306. _STARPU_MPI_DEBUG(1, "Sending data %p to %d\n", data_handle, node);
  307. starpu_mpi_isend_detached(data_handle, node, tag, comm, NULL, NULL);
  308. }
  309. }
  310. }
  311. void starpu_mpi_get_data_on_node(MPI_Comm comm, starpu_data_handle_t data_handle, int node)
  312. {
  313. int me, rank, tag;
  314. rank = starpu_mpi_data_get_rank(data_handle);
  315. if (rank == -1)
  316. {
  317. _STARPU_ERROR("StarPU needs to be told the MPI rank of this data, using starpu_mpi_data_register\n");
  318. }
  319. starpu_mpi_comm_rank(comm, &me);
  320. if (node == rank)
  321. return;
  322. tag = starpu_mpi_data_get_tag(data_handle);
  323. if (tag == -1)
  324. {
  325. _STARPU_ERROR("StarPU needs to be told the MPI tag of this data, using starpu_mpi_data_register\n");
  326. }
  327. if (me == node)
  328. {
  329. MPI_Status status;
  330. _STARPU_MPI_DEBUG(1, "Migrating data %p from %d to %d\n", data_handle, rank, node);
  331. int already_received = _starpu_mpi_cache_received_data_set(data_handle);
  332. if (already_received == 0)
  333. {
  334. _STARPU_MPI_DEBUG(1, "Receiving data %p from %d\n", data_handle, rank);
  335. starpu_mpi_recv(data_handle, rank, tag, comm, &status);
  336. }
  337. }
  338. else if (me == rank)
  339. {
  340. _STARPU_MPI_DEBUG(1, "Migrating data %p from %d to %d\n", data_handle, rank, node);
  341. int already_sent = _starpu_mpi_cache_sent_data_set(data_handle, node);
  342. if (already_sent == 0)
  343. {
  344. _STARPU_MPI_DEBUG(1, "Sending data %p to %d\n", data_handle, node);
  345. starpu_mpi_send(data_handle, node, tag, comm);
  346. }
  347. }
  348. }
  349. void starpu_mpi_get_data_on_all_nodes_detached(MPI_Comm comm, starpu_data_handle_t data_handle)
  350. {
  351. int size, i;
  352. starpu_mpi_comm_size(comm, &size);
  353. for (i = 0; i < size; i++)
  354. starpu_mpi_get_data_on_node_detached(comm, data_handle, i, NULL, NULL);
  355. }
  356. void starpu_mpi_data_migrate(MPI_Comm comm, starpu_data_handle_t data, int new_rank)
  357. {
  358. int old_rank = starpu_mpi_data_get_rank(data);
  359. if (new_rank == old_rank)
  360. /* Already there */
  361. return;
  362. /* First submit data migration if it's not already on destination */
  363. starpu_mpi_get_data_on_node_detached(comm, data, new_rank, NULL, NULL);
  364. /* And note new owner */
  365. starpu_mpi_data_set_rank_comm(data, new_rank, comm);
  366. /* Flush cache in all other nodes */
  367. /* TODO: Ideally we'd transmit the knowledge of who owns it */
  368. starpu_mpi_cache_flush(comm, data);
  369. return;
  370. }
  371. int starpu_mpi_wait_for_all(MPI_Comm comm)
  372. {
  373. int mpi = 1;
  374. int task = 1;
  375. while (task || mpi)
  376. {
  377. task = _starpu_task_wait_for_all_and_return_nb_waited_tasks();
  378. mpi = _starpu_mpi_barrier(comm);
  379. }
  380. return 0;
  381. }