deal.II version GIT relicensing-6834-g5b78e6bcdf 2026-10-01 11:20:01+00:00
\(\newcommand{\dealvcentcolon}{\mathrel{\mathop{:}}}\) \(\newcommand{\dealcoloneq}{\dealvcentcolon\mathrel{\mkern-1.2mu}=}\) \(\newcommand{\jump}[1]{\left[\!\left[ #1 \right]\!\right]}\) \(\newcommand{\average}[1]{\left\{\!\left\{ #1 \right\}\!\right\}}\)
Loading...
Searching...
No Matches
mpi_consensus_algorithms.h
Go to the documentation of this file.
1// -----------------------------------------------------------------------------
2//
3// SPDX-License-Identifier: Apache-2.0 WITH LLVM-exception OR LGPL-2.1-or-later
4// Copyright (C) 2020 - 2026 by the deal.II authors
5//
6// This file is part of the deal.II library.
7//
8// Detailed license information governing the source code and contributions
9// can be found in LICENSE.md and CONTRIBUTING.md at the top level directory.
10//
11// -----------------------------------------------------------------------------
12
13#ifndef dealii_mpi_consensus_algorithm_h
14#define dealii_mpi_consensus_algorithm_h
15
16#include <deal.II/base/config.h>
17
18#include <deal.II/base/mpi.h>
19#include <deal.II/base/mpi.templates.h>
21
22#include <set>
23#include <vector>
24
26
27
28namespace Utilities
29{
30 namespace MPI
31 {
130 namespace ConsensusAlgorithms
131 {
145 template <typename RequestType, typename AnswerType>
147 {
148 public:
152 Interface() = default;
153
158 virtual ~Interface() = default;
159
183 virtual std::vector<unsigned int>
185 const std::vector<unsigned int> &targets,
186 const std::function<RequestType(const unsigned int)> &create_request,
187 const std::function<AnswerType(const unsigned int,
188 const RequestType &)> &answer_request,
189 const std::function<void(const unsigned int, const AnswerType &)>
190 &process_answer,
191 const MPI_Comm comm) = 0;
192 };
193
194
208 template <typename RequestType, typename AnswerType>
209 class NBX : public Interface<RequestType, AnswerType>
210 {
211 public:
215 NBX() = default;
216
220 virtual ~NBX() = default;
221
225 virtual std::vector<unsigned int>
227 const std::vector<unsigned int> &targets,
228 const std::function<RequestType(const unsigned int)> &create_request,
229 const std::function<AnswerType(const unsigned int,
230 const RequestType &)> &answer_request,
231 const std::function<void(const unsigned int, const AnswerType &)>
232 &process_answer,
233 const MPI_Comm comm) override;
234
235 private:
236#ifdef DEAL_II_WITH_MPI
240 std::vector<std::vector<char>> send_buffers;
241
245 std::vector<MPI_Request> send_requests;
246
254 std::vector<std::unique_ptr<std::vector<char>>> request_buffers;
255
259 std::vector<std::unique_ptr<MPI_Request>> request_requests;
260
265
266 // request for barrier
267 MPI_Request barrier_request;
268#endif
269
273 std::set<unsigned int> requesting_processes;
274
280 bool
282 const std::function<void(const unsigned int, const AnswerType &)>
283 &process_answer,
284 const MPI_Comm comm);
285
290 void
292
298 bool
300
306 void
308 const std::function<AnswerType(const unsigned int,
309 const RequestType &)> &answer_request,
310 const MPI_Comm comm);
311
316 void
318 const std::vector<unsigned int> &targets,
319 const std::function<RequestType(const unsigned int)> &create_request,
320 const MPI_Comm comm);
321
326 void
328 };
329
330
375 template <typename RequestType, typename AnswerType>
376 std::vector<unsigned int>
377 nbx(const std::vector<unsigned int> &targets,
378 const std::function<RequestType(const unsigned int)> &create_request,
379 const std::function<AnswerType(const unsigned int,
380 const RequestType &)> &answer_request,
381 const std::function<void(const unsigned int, const AnswerType &)>
382 &process_answer,
383 const MPI_Comm comm);
384
422 template <typename RequestType>
423 std::vector<unsigned int>
424 nbx(const std::vector<unsigned int> &targets,
425 const std::function<RequestType(const unsigned int)> &create_request,
426 const std::function<void(const unsigned int, const RequestType &)>
427 &process_request,
428 const MPI_Comm comm);
429
455 template <typename RequestType, typename AnswerType>
456 class PEX : public Interface<RequestType, AnswerType>
457 {
458 public:
462 PEX() = default;
463
467 virtual ~PEX() = default;
468
472 virtual std::vector<unsigned int>
474 const std::vector<unsigned int> &targets,
475 const std::function<RequestType(const unsigned int)> &create_request,
476 const std::function<AnswerType(const unsigned int,
477 const RequestType &)> &answer_request,
478 const std::function<void(const unsigned int, const AnswerType &)>
479 &process_answer,
480 const MPI_Comm comm) override;
481
482 private:
483#ifdef DEAL_II_WITH_MPI
487 std::vector<std::vector<char>> send_buffers;
488
492 std::vector<std::vector<char>> recv_buffers;
493
497 std::vector<MPI_Request> send_request_requests;
498
502 std::vector<std::vector<char>> requests_buffers;
503
507 std::vector<MPI_Request> send_answer_requests;
508#endif
512 std::set<unsigned int> requesting_processes;
513
518 unsigned int
520 const std::vector<unsigned int> &targets,
521 const std::function<RequestType(const unsigned int)> &create_request,
522 const MPI_Comm comm);
523
528 void
530 const unsigned int index,
531 const std::function<AnswerType(const unsigned int,
532 const RequestType &)> &answer_request,
533 const MPI_Comm comm);
534
539 void
541 const unsigned int n_targets,
542 const std::function<void(const unsigned int, const AnswerType &)>
543 &process_answer,
544 const MPI_Comm comm);
545
550 void
552 };
553
554
555
612 template <typename RequestType, typename AnswerType>
613 std::vector<unsigned int>
614 pex(const std::vector<unsigned int> &targets,
615 const std::function<RequestType(const unsigned int)> &create_request,
616 const std::function<AnswerType(const unsigned int,
617 const RequestType &)> &answer_request,
618 const std::function<void(const unsigned int, const AnswerType &)>
619 &process_answer,
620 const MPI_Comm comm);
621
659 template <typename RequestType>
660 std::vector<unsigned int>
661 pex(const std::vector<unsigned int> &targets,
662 const std::function<RequestType(const unsigned int)> &create_request,
663 const std::function<void(const unsigned int, const RequestType &)>
664 &process_request,
665 const MPI_Comm comm);
666
667
672 template <typename RequestType, typename AnswerType>
673 class Serial : public Interface<RequestType, AnswerType>
674 {
675 public:
679 Serial() = default;
680
684 virtual std::vector<unsigned int>
686 const std::vector<unsigned int> &targets,
687 const std::function<RequestType(const unsigned int)> &create_request,
688 const std::function<AnswerType(const unsigned int,
689 const RequestType &)> &answer_request,
690 const std::function<void(const unsigned int, const AnswerType &)>
691 &process_answer,
692 const MPI_Comm comm) override;
693 };
694
695
696
729 template <typename RequestType, typename AnswerType>
730 std::vector<unsigned int>
732 const std::vector<unsigned int> &targets,
733 const std::function<RequestType(const unsigned int)> &create_request,
734 const std::function<AnswerType(const unsigned int, const RequestType &)>
735 &answer_request,
736 const std::function<void(const unsigned int, const AnswerType &)>
737 &process_answer,
738 const MPI_Comm comm);
739
769 template <typename RequestType>
770 std::vector<unsigned int>
772 const std::vector<unsigned int> &targets,
773 const std::function<RequestType(const unsigned int)> &create_request,
774 const std::function<void(const unsigned int, const RequestType &)>
775 &process_request,
776 const MPI_Comm comm);
777
778
779
792 template <typename RequestType, typename AnswerType>
793 class Selector : public Interface<RequestType, AnswerType>
794 {
795 public:
799 Selector() = default;
800
804 virtual ~Selector() = default;
805
811 virtual std::vector<unsigned int>
813 const std::vector<unsigned int> &targets,
814 const std::function<RequestType(const unsigned int)> &create_request,
815 const std::function<AnswerType(const unsigned int,
816 const RequestType &)> &answer_request,
817 const std::function<void(const unsigned int, const AnswerType &)>
818 &process_answer,
819 const MPI_Comm comm) override;
820
821 private:
822 // Pointer to the actual ConsensusAlgorithms::Interface implementation.
823 std::shared_ptr<Interface<RequestType, AnswerType>> consensus_algo;
824 };
825
826
827
872 template <typename RequestType, typename AnswerType>
873 std::vector<unsigned int>
875 const std::vector<unsigned int> &targets,
876 const std::function<RequestType(const unsigned int)> &create_request,
877 const std::function<AnswerType(const unsigned int, const RequestType &)>
878 &answer_request,
879 const std::function<void(const unsigned int, const AnswerType &)>
880 &process_answer,
881 const MPI_Comm comm);
882
920 template <typename RequestType>
921 std::vector<unsigned int>
923 const std::vector<unsigned int> &targets,
924 const std::function<RequestType(const unsigned int)> &create_request,
925 const std::function<void(const unsigned int, const RequestType &)>
926 &process_request,
927 const MPI_Comm comm);
928
929
930
931#ifndef DOXYGEN
932 // Implementation of the functions in this namespace.
933
934 template <typename RequestType, typename AnswerType>
935 std::vector<unsigned int>
936 nbx(const std::vector<unsigned int> &targets,
937 const std::function<RequestType(const unsigned int)> &create_request,
938 const std::function<AnswerType(const unsigned int,
939 const RequestType &)> &answer_request,
940 const std::function<void(const unsigned int, const AnswerType &)>
941 &process_answer,
942 const MPI_Comm comm)
943 {
945 targets, create_request, answer_request, process_answer, comm);
946 }
947
948
949
950 template <typename RequestType>
951 std::vector<unsigned int>
952 nbx(const std::vector<unsigned int> &targets,
953 const std::function<RequestType(const unsigned int)> &create_request,
954 const std::function<void(const unsigned int, const RequestType &)>
955 &process_request,
956 const MPI_Comm comm)
957 {
958 // TODO: For the moment, simply implement this special case by
959 // forwarding to the other function with rewritten function
960 // objects and using an empty type as answer type. This way,
961 // we have the interface in place and can provide a more
962 // efficient implementation later on.
963 using EmptyType = std::tuple<>;
964
965 return nbx<RequestType, EmptyType>(
966 targets,
967 create_request,
968 // answer_request:
969 [&process_request](const unsigned int source_rank,
970 const RequestType &request) -> EmptyType {
971 process_request(source_rank, request);
972 // Return something. What it is is arbitrary here, except that
973 // we want it to be as small an object as possible. Using
974 // std::tuple<> is interpreted as an empty object that is packed
975 // down to a zero-length char array.
976 return {};
977 },
978 // process_answer:
979 [](const unsigned int /*target_rank */,
980 const EmptyType & /*answer*/) {},
981 comm);
982 }
983
984
985
986 template <typename RequestType, typename AnswerType>
987 std::vector<unsigned int>
988 pex(const std::vector<unsigned int> &targets,
989 const std::function<RequestType(const unsigned int)> &create_request,
990 const std::function<AnswerType(const unsigned int,
991 const RequestType &)> &answer_request,
992 const std::function<void(const unsigned int, const AnswerType &)>
993 &process_answer,
994 const MPI_Comm comm)
995 {
996 return PEX<RequestType, AnswerType>().run(
997 targets, create_request, answer_request, process_answer, comm);
998 }
999
1000
1001
1002 template <typename RequestType>
1003 std::vector<unsigned int>
1004 pex(const std::vector<unsigned int> &targets,
1005 const std::function<RequestType(const unsigned int)> &create_request,
1006 const std::function<void(const unsigned int, const RequestType &)>
1007 &process_request,
1008 const MPI_Comm comm)
1009 {
1010 // TODO: For the moment, simply implement this special case by
1011 // forwarding to the other function with rewritten function
1012 // objects and using an empty type as answer type. This way,
1013 // we have the interface in place and can provide a more
1014 // efficient implementation later on.
1015 using EmptyType = std::tuple<>;
1016
1017 return pex<RequestType, EmptyType>(
1018 targets,
1019 create_request,
1020 // answer_request:
1021 [&process_request](const unsigned int source_rank,
1022 const RequestType &request) -> EmptyType {
1023 process_request(source_rank, request);
1024 // Return something. What it is is arbitrary here, except that
1025 // we want it to be as small an object as possible. Using
1026 // std::tuple<> is interpreted as an empty object that is packed
1027 // down to a zero-length char array.
1028 return {};
1029 },
1030 // process_answer:
1031 [](const unsigned int /*target_rank */,
1032 const EmptyType & /*answer*/) {},
1033 comm);
1034 }
1035
1036
1037
1038 template <typename RequestType, typename AnswerType>
1039 std::vector<unsigned int>
1040 serial(
1041 const std::vector<unsigned int> &targets,
1042 const std::function<RequestType(const unsigned int)> &create_request,
1043 const std::function<AnswerType(const unsigned int, const RequestType &)>
1044 &answer_request,
1045 const std::function<void(const unsigned int, const AnswerType &)>
1046 &process_answer,
1047 const MPI_Comm comm)
1048 {
1049 return Serial<RequestType, AnswerType>().run(
1050 targets, create_request, answer_request, process_answer, comm);
1051 }
1052
1053
1054
1055 template <typename RequestType>
1056 std::vector<unsigned int>
1057 serial(
1058 const std::vector<unsigned int> &targets,
1059 const std::function<RequestType(const unsigned int)> &create_request,
1060 const std::function<void(const unsigned int, const RequestType &)>
1061 &process_request,
1062 const MPI_Comm comm)
1063 {
1064 // TODO: For the moment, simply implement this special case by
1065 // forwarding to the other function with rewritten function
1066 // objects and using an empty type as answer type. This way,
1067 // we have the interface in place and can provide a more
1068 // efficient implementation later on.
1069 using EmptyType = std::tuple<>;
1070
1071 return serial<RequestType, EmptyType>(
1072 targets,
1073 create_request,
1074 // answer_request:
1075 [&process_request](const unsigned int source_rank,
1076 const RequestType &request) -> EmptyType {
1077 process_request(source_rank, request);
1078 // Return something. What it is is arbitrary here, except that
1079 // we want it to be as small an object as possible. Using
1080 // std::tuple<> is interpreted as an empty object that is packed
1081 // down to a zero-length char array.
1082 return {};
1083 },
1084 // process_answer:
1085 [](const unsigned int /*target_rank */,
1086 const EmptyType & /*answer*/) {},
1087 comm);
1088 }
1089
1090
1091
1092 template <typename RequestType, typename AnswerType>
1093 std::vector<unsigned int>
1094 selector(
1095 const std::vector<unsigned int> &targets,
1096 const std::function<RequestType(const unsigned int)> &create_request,
1097 const std::function<AnswerType(const unsigned int, const RequestType &)>
1098 &answer_request,
1099 const std::function<void(const unsigned int, const AnswerType &)>
1100 &process_answer,
1101 const MPI_Comm comm)
1102 {
1103 return Selector<RequestType, AnswerType>().run(
1104 targets, create_request, answer_request, process_answer, comm);
1105 }
1106
1107
1108
1109 template <typename RequestType>
1110 std::vector<unsigned int>
1111 selector(
1112 const std::vector<unsigned int> &targets,
1113 const std::function<RequestType(const unsigned int)> &create_request,
1114 const std::function<void(const unsigned int, const RequestType &)>
1115 &process_request,
1116 const MPI_Comm comm)
1117 {
1118 // TODO: For the moment, simply implement this special case by
1119 // forwarding to the other function with rewritten function
1120 // objects and using an empty type as answer type. This way,
1121 // we have the interface in place and can provide a more
1122 // efficient implementation later on.
1123 using EmptyType = std::tuple<>;
1124
1125 return selector<RequestType, EmptyType>(
1126 targets,
1127 create_request,
1128 // answer_request:
1129 [&process_request](const unsigned int source_rank,
1130 const RequestType &request) -> EmptyType {
1131 process_request(source_rank, request);
1132 // Return something. What it is is arbitrary here, except that
1133 // we want it to be as small an object as possible. Using
1134 // std::tuple<> is interpreted as an empty object that is packed
1135 // down to a zero-length char array.
1136 return {};
1137 },
1138 // process_answer:
1139 [](const unsigned int /*target_rank */,
1140 const EmptyType & /*answer*/) {},
1141 comm);
1142 }
1143
1144#endif
1145
1146
1147 } // namespace ConsensusAlgorithms
1148 } // end of namespace MPI
1149} // end of namespace Utilities
1150
1151
1152
1153#ifndef DOXYGEN
1154
1155// ----------------- Implementation of template functions
1156
1157namespace Utilities
1158{
1159 namespace MPI
1160 {
1161 namespace ConsensusAlgorithms
1162 {
1163 namespace internal
1164 {
1169 inline bool
1170 has_unique_elements(const std::vector<unsigned int> &targets)
1171 {
1172 std::vector<unsigned int> my_destinations = targets;
1173 std::sort(my_destinations.begin(), my_destinations.end());
1174 return (std::adjacent_find(my_destinations.begin(),
1175 my_destinations.end()) ==
1176 my_destinations.end());
1177 }
1178
1179
1180
1184 inline void
1185 handle_exception(std::exception_ptr &&exception, const MPI_Comm comm)
1186 {
1187# ifdef DEAL_II_WITH_MPI
1188 // an exception within a ConsensusAlgorithm likely causes an
1189 // MPI deadlock. Abort with a reasonable error message instead.
1190 try
1191 {
1192 std::rethrow_exception(exception);
1193 }
1194 catch (ExceptionBase &exc)
1195 {
1196 // report name of the deal.II exception:
1197 std::cerr
1198 << std::endl
1199 << std::endl
1200 << "----------------------------------------------------"
1201 << std::endl;
1202 std::cerr
1203 << "Exception '" << exc.get_exc_name() << "'"
1204 << " on rank " << Utilities::MPI::this_mpi_process(comm)
1205 << " on processing: " << std::endl
1206 << exc.what() << std::endl
1207 << "Aborting!" << std::endl
1208 << "----------------------------------------------------"
1209 << std::endl;
1210
1211 // Then bring down the whole MPI world
1212 MPI_Abort(comm, 255);
1213 }
1214 catch (std::exception &exc)
1215 {
1216 std::cerr
1217 << std::endl
1218 << std::endl
1219 << "----------------------------------------------------"
1220 << std::endl;
1221 std::cerr
1222 << "Exception within ConsensusAlgorithm"
1223 << " on rank " << Utilities::MPI::this_mpi_process(comm)
1224 << " on processing: " << std::endl
1225 << exc.what() << std::endl
1226 << "Aborting!" << std::endl
1227 << "----------------------------------------------------"
1228 << std::endl;
1229
1230 // Then bring down the whole MPI world
1231 MPI_Abort(comm, 255);
1232 }
1233 catch (...)
1234 {
1235 std::cerr
1236 << std::endl
1237 << std::endl
1238 << "----------------------------------------------------"
1239 << std::endl;
1240 std::cerr
1241 << "Unknown exception within ConsensusAlgorithm!" << std::endl
1242 << "Aborting!" << std::endl
1243 << "----------------------------------------------------"
1244 << std::endl;
1245
1246 // Then bring down the whole MPI world
1247 MPI_Abort(comm, 255);
1248 }
1249# else
1250 (void)comm;
1251
1252 // No need to be concerned about deadlocks without MPI.
1253 // Defer to exception handling further up the callstack.
1254 std::rethrow_exception(exception);
1255# endif
1256 }
1257 } // namespace internal
1258
1259
1260
1261 template <typename RequestType, typename AnswerType>
1262 std::vector<unsigned int>
1264 const std::vector<unsigned int> &targets,
1265 const std::function<RequestType(const unsigned int)> &create_request,
1266 const std::function<AnswerType(const unsigned int, const RequestType &)>
1267 &answer_request,
1268 const std::function<void(const unsigned int, const AnswerType &)>
1269 &process_answer,
1270 const MPI_Comm comm)
1271 {
1272 Assert(internal::has_unique_elements(targets),
1273 ExcMessage("The consensus algorithms expect that each process "
1274 "only sends a single message to another process, "
1275 "but the targets provided include duplicates."));
1276
1277 static CollectiveMutex mutex;
1278 CollectiveMutex::ScopedLock lock(mutex, comm);
1279
1280 try
1281 {
1282 // 1) Send data to identified targets and start receiving
1283 // the answers from these very same processes.
1284 start_communication(targets, create_request, comm);
1285
1286 // 2) Until all posted receive operations are known to have
1287 // completed, answer requests and keep checking whether all
1288 // requests of this process have been answered.
1289 //
1290 // The requests that we catch in the answer_requests()
1291 // function originate elsewhere, that is, they are not in
1292 // response to our own messages
1293 //
1294 // Note also that we may not catch all incoming requests in
1295 // the following two lines: our own requests may have been
1296 // satisfied before we've dealt with all incoming requests.
1297 // That's ok: We will get around to dealing with all
1298 // remaining message later. We just want to move on to the
1299 // next step as early as possible.
1300 while (all_locally_originated_receives_are_completed(process_answer,
1301 comm) == false)
1302 maybe_answer_one_request(answer_request, comm);
1303
1304 // 3) Signal to all other processes that all requests of this
1305 // process have been answered
1306 signal_finish(comm);
1307
1308 // 4) Nevertheless, this process has to keep on answering
1309 // (potential) incoming requests until all processes have
1310 // received the answer to all requests
1311 while (all_remotely_originated_receives_are_completed() == false)
1312 maybe_answer_one_request(answer_request, comm);
1313
1314 // 5) process the answer to all requests
1315 clean_up_and_end_communication(comm);
1316 }
1317 catch (...)
1318 {
1319 internal::handle_exception(std::current_exception(), comm);
1320 }
1321
1322 return std::vector<unsigned int>(requesting_processes.begin(),
1323 requesting_processes.end());
1324 }
1325
1326
1327
1328 template <typename RequestType, typename AnswerType>
1329 void
1331 const std::vector<unsigned int> &targets,
1332 const std::function<RequestType(const unsigned int)> &create_request,
1333 const MPI_Comm comm)
1334 {
1335# ifdef DEAL_II_WITH_MPI
1336 // 1)
1337 const auto n_targets = targets.size();
1338
1339 const int tag_request = Utilities::MPI::internal::Tags::
1341
1342 // 2) allocate memory
1343 send_requests.resize(n_targets);
1344 send_buffers.resize(n_targets);
1345
1346 {
1347 // 4) send and receive
1348 for (unsigned int index = 0; index < n_targets; ++index)
1349 {
1350 const unsigned int rank = targets[index];
1352
1353 auto &send_buffer = send_buffers[index];
1354 send_buffer =
1355 (create_request ? Utilities::pack(create_request(rank), false) :
1356 std::vector<char>());
1357
1358 // Post a request to send data
1359 auto ierr = MPI_Isend(send_buffer.data(),
1360 send_buffer.size(),
1361 MPI_CHAR,
1362 rank,
1363 tag_request,
1364 comm,
1365 &send_requests[index]);
1366 AssertThrowMPI(ierr);
1367 }
1368
1369 // Also record that we expect an answer from each target we sent
1370 // a request to:
1371 n_outstanding_answers = n_targets;
1372 }
1373# else
1374 (void)targets;
1375 (void)create_request;
1376 (void)comm;
1377# endif
1378 }
1379
1380
1381
1382 template <typename RequestType, typename AnswerType>
1383 bool
1386 const std::function<void(const unsigned int, const AnswerType &)>
1387 &process_answer,
1388 const MPI_Comm comm)
1389 {
1390# ifdef DEAL_II_WITH_MPI
1391 // We know that all requests have come in when we have pending
1392 // messages from all targets with the right tag (some of which we may
1393 // have already taken care of below, after discovering their existence).
1394 // We can check for pending messages with MPI_IProbe, which returns
1395 // immediately with a return code that indicates whether
1396 // it has found a message from any process with a given
1397 // tag.
1398 if (n_outstanding_answers == 0)
1399 return true;
1400 else
1401 {
1402 const int tag_deliver = Utilities::MPI::internal::Tags::
1404
1405 int request_is_pending;
1406 MPI_Status status;
1407 const auto ierr = MPI_Iprobe(
1408 MPI_ANY_SOURCE, tag_deliver, comm, &request_is_pending, &status);
1409 AssertThrowMPI(ierr);
1410
1411 // If there is no pending message with this tag,
1412 // then we are clearly not done receiving everything
1413 // yet -- so return false.
1414 if (request_is_pending == 0)
1415 return false;
1416 else
1417 {
1418 // OK, so we have gotten a reply to our request from
1419 // one rank. Let us process it.
1420 const auto target = status.MPI_SOURCE;
1421
1422 // Then query the size of the message, allocate enough memory,
1423 // receive the data, and process it.
1424 int message_size;
1425 {
1426 const int ierr =
1427 MPI_Get_count(&status, MPI_CHAR, &message_size);
1428 AssertThrowMPI(ierr);
1429 }
1430 std::vector<char> recv_buffer(message_size);
1431
1432 {
1433 const int tag_deliver = Utilities::MPI::internal::Tags::
1435
1436 const int ierr = MPI_Recv(recv_buffer.data(),
1437 recv_buffer.size(),
1438 MPI_CHAR,
1439 target,
1440 tag_deliver,
1441 comm,
1442 MPI_STATUS_IGNORE);
1443 AssertThrowMPI(ierr);
1444 }
1445
1446 if (process_answer)
1447 process_answer(target,
1448 Utilities::unpack<AnswerType>(recv_buffer,
1449 false));
1450
1451 // Finally, remove this rank from the list of outstanding
1452 // targets:
1453 --n_outstanding_answers;
1454
1455 // We could do another go-around from the top of this
1456 // else-branch to see whether there are actually other messages
1457 // that are currently pending. But that would mean spending
1458 // substantial time in receiving answers while we should also be
1459 // sending answers to requests we have received from other
1460 // places. So let it be enough for now. If there are outstanding
1461 // answers, we will get back to this function before long and
1462 // can take care of them then.
1463 return (n_outstanding_answers == 0);
1464 }
1465 }
1466
1467# else
1468 (void)process_answer;
1469 (void)comm;
1470
1471 return true;
1472# endif
1473 }
1474
1475
1476
1477 template <typename RequestType, typename AnswerType>
1478 void
1480 const std::function<AnswerType(const unsigned int, const RequestType &)>
1481 &answer_request,
1482 const MPI_Comm comm)
1483 {
1484# ifdef DEAL_II_WITH_MPI
1485
1486 const int tag_request = Utilities::MPI::internal::Tags::
1488 const int tag_deliver = Utilities::MPI::internal::Tags::
1490
1491 // Check if there is a request pending. By selecting the
1492 // tag_request tag, these are other processes asking for
1493 // our own replies, not these other processes' replies
1494 // to our own requests.
1495 //
1496 // There may be multiple such pending messages. We
1497 // only answer one.
1498 MPI_Status status;
1499 int request_is_pending;
1500 const auto ierr = MPI_Iprobe(
1501 MPI_ANY_SOURCE, tag_request, comm, &request_is_pending, &status);
1502 AssertThrowMPI(ierr);
1503
1504 if (request_is_pending != 0)
1505 {
1506 // Get the rank of the requesting process and add it to the
1507 // list of requesting processes (which may contain duplicates).
1508 const auto other_rank = status.MPI_SOURCE;
1509
1510 Assert(requesting_processes.find(other_rank) ==
1511 requesting_processes.end(),
1512 ExcMessage("Process is requesting a second time!"));
1513 requesting_processes.insert(other_rank);
1514
1515 // get size of incoming message
1516 int number_amount;
1517 auto ierr = MPI_Get_count(&status, MPI_CHAR, &number_amount);
1518 AssertThrowMPI(ierr);
1519
1520 // allocate memory for incoming message
1521 std::vector<char> buffer_recv(number_amount);
1522 ierr = MPI_Recv(buffer_recv.data(),
1523 number_amount,
1524 MPI_CHAR,
1525 other_rank,
1526 tag_request,
1527 comm,
1528 MPI_STATUS_IGNORE);
1529 AssertThrowMPI(ierr);
1530
1531 // Allocate memory for an answer message to the current request,
1532 // and ask the 'process' object to produce an answer:
1533 request_buffers.emplace_back(std::make_unique<std::vector<char>>());
1534 auto &request_buffer = *request_buffers.back();
1535 if (answer_request)
1536 request_buffer =
1537 Utilities::pack(answer_request(other_rank,
1538 Utilities::unpack<RequestType>(
1539 buffer_recv, false)),
1540 false);
1541
1542 // Then initiate sending the answer back to the requester.
1543 request_requests.emplace_back(std::make_unique<MPI_Request>());
1544 ierr = MPI_Isend(request_buffer.data(),
1545 request_buffer.size(),
1546 MPI_CHAR,
1547 other_rank,
1548 tag_deliver,
1549 comm,
1550 request_requests.back().get());
1551 AssertThrowMPI(ierr);
1552 }
1553# else
1554 (void)answer_request;
1555 (void)comm;
1556# endif
1557 }
1558
1559
1560
1561 template <typename RequestType, typename AnswerType>
1562 void
1564 {
1565# ifdef DEAL_II_WITH_MPI
1566 const auto ierr = MPI_Ibarrier(comm, &barrier_request);
1567 AssertThrowMPI(ierr);
1568# else
1569 (void)comm;
1570# endif
1571 }
1572
1573
1574
1575 template <typename RequestType, typename AnswerType>
1576 bool
1577 NBX<RequestType,
1578 AnswerType>::all_remotely_originated_receives_are_completed()
1579 {
1580# ifdef DEAL_II_WITH_MPI
1581 int all_ranks_reached_barrier;
1582 const auto ierr = MPI_Test(&barrier_request,
1583 &all_ranks_reached_barrier,
1584 MPI_STATUS_IGNORE);
1585 AssertThrowMPI(ierr);
1586 return all_ranks_reached_barrier != 0;
1587# else
1588 return true;
1589# endif
1590 }
1591
1592
1593
1594 template <typename RequestType, typename AnswerType>
1595 void
1597 const MPI_Comm comm)
1598 {
1599 (void)comm;
1600# ifdef DEAL_II_WITH_MPI
1601 // clean up
1602 {
1603 if (send_requests.size() > 0)
1604 {
1605 const int ierr = MPI_Waitall(send_requests.size(),
1606 send_requests.data(),
1607 MPI_STATUSES_IGNORE);
1608 AssertThrowMPI(ierr);
1609 }
1610
1611 int ierr = MPI_Wait(&barrier_request, MPI_STATUS_IGNORE);
1612 AssertThrowMPI(ierr);
1613
1614 for (auto &i : request_requests)
1615 {
1616 ierr = MPI_Wait(i.get(), MPI_STATUS_IGNORE);
1617 AssertThrowMPI(ierr);
1618 }
1619
1620 if constexpr (running_in_debug_mode())
1621 {
1622 // note: IBarrier seems to make problem during testing, this
1623 // additional Barrier seems to help
1624 ierr = MPI_Barrier(comm);
1625 AssertThrowMPI(ierr);
1626 }
1627 }
1628# endif
1629 }
1630
1631
1632
1633 template <typename RequestType, typename AnswerType>
1634 std::vector<unsigned int>
1636 const std::vector<unsigned int> &targets,
1637 const std::function<RequestType(const unsigned int)> &create_request,
1638 const std::function<AnswerType(const unsigned int, const RequestType &)>
1639 &answer_request,
1640 const std::function<void(const unsigned int, const AnswerType &)>
1641 &process_answer,
1642 const MPI_Comm comm)
1643 {
1644 Assert(internal::has_unique_elements(targets),
1645 ExcMessage("The consensus algorithms expect that each process "
1646 "only sends a single message to another process, "
1647 "but the targets provided include duplicates."));
1648
1649 static CollectiveMutex mutex;
1650 CollectiveMutex::ScopedLock lock(mutex, comm);
1651
1652 try
1653 {
1654 // 1) Send requests and start receiving the answers.
1655 // In particular, determine how many requests we should expect
1656 // on the current process.
1657 const unsigned int n_requests =
1658 start_communication(targets, create_request, comm);
1659
1660 // 2) Answer requests:
1661 for (unsigned int request = 0; request < n_requests; ++request)
1662 answer_one_request(request, answer_request, comm);
1663
1664 // 3) Process answers:
1665 process_incoming_answers(targets.size(), process_answer, comm);
1666
1667 // 4) Make sure all sends have successfully terminated:
1668 clean_up_and_end_communication();
1669 }
1670 catch (...)
1671 {
1672 internal::handle_exception(std::current_exception(), comm);
1673 }
1674
1675 return std::vector<unsigned int>(requesting_processes.begin(),
1676 requesting_processes.end());
1677 }
1678
1679
1680
1681 template <typename RequestType, typename AnswerType>
1682 unsigned int
1684 const std::vector<unsigned int> &targets,
1685 const std::function<RequestType(const unsigned int)> &create_request,
1686 const MPI_Comm comm)
1687 {
1688# ifdef DEAL_II_WITH_MPI
1689 const int tag_request = Utilities::MPI::internal::Tags::
1691
1692 // 1) determine with which processes this process wants to communicate
1693 // with
1694 const unsigned int n_targets = targets.size();
1695
1696 // 2) determine who wants to communicate with this process
1697 const unsigned int n_sources =
1699
1700 // 2) allocate memory
1701 recv_buffers.resize(n_targets);
1702 send_buffers.resize(n_targets);
1703 send_request_requests.resize(n_targets);
1704
1705 send_answer_requests.resize(n_sources);
1706 requests_buffers.resize(n_sources);
1707
1708 // 4) send and receive
1709 for (unsigned int i = 0; i < n_targets; ++i)
1710 {
1711 const unsigned int rank = targets[i];
1713
1714 // pack data which should be sent
1715 auto &send_buffer = send_buffers[i];
1716 if (create_request)
1717 send_buffer = Utilities::pack(create_request(rank), false);
1718
1719 // start to send data
1720 auto ierr = MPI_Isend(send_buffer.data(),
1721 send_buffer.size(),
1722 MPI_CHAR,
1723 rank,
1724 tag_request,
1725 comm,
1726 &send_request_requests[i]);
1727 AssertThrowMPI(ierr);
1728 }
1729
1730 return n_sources;
1731# else
1732 (void)targets;
1733 (void)create_request;
1734 (void)comm;
1735 return 0;
1736# endif
1737 }
1738
1739
1740
1741 template <typename RequestType, typename AnswerType>
1742 void
1744 const unsigned int index,
1745 const std::function<AnswerType(const unsigned int, const RequestType &)>
1746 &answer_request,
1747 const MPI_Comm comm)
1748 {
1749# ifdef DEAL_II_WITH_MPI
1750 const int tag_request = Utilities::MPI::internal::Tags::
1752 const int tag_deliver = Utilities::MPI::internal::Tags::
1754
1755 // Wait until we have a message ready for retrieval, though we don't
1756 // care which process it is from.
1757 MPI_Status status;
1758 int ierr = MPI_Probe(MPI_ANY_SOURCE, tag_request, comm, &status);
1759 AssertThrowMPI(ierr);
1760
1761 // Get rank of incoming message and verify that it makes sense
1762 const unsigned int other_rank = status.MPI_SOURCE;
1763
1764 Assert(requesting_processes.find(other_rank) ==
1765 requesting_processes.end(),
1766 ExcMessage(
1767 "A process is sending a request after a request from "
1768 "the same process has previously already been "
1769 "received. This algorithm does not expect this to happen."));
1770 requesting_processes.insert(other_rank);
1771
1772 // Actually get the incoming message:
1773 int number_amount;
1774 ierr = MPI_Get_count(&status, MPI_CHAR, &number_amount);
1775 AssertThrowMPI(ierr);
1776
1777 std::vector<char> buffer_recv(number_amount);
1778 ierr = MPI_Recv(buffer_recv.data(),
1779 number_amount,
1780 MPI_CHAR,
1781 other_rank,
1782 tag_request,
1783 comm,
1784 &status);
1785 AssertThrowMPI(ierr);
1786
1787 // Process request by asking the user-provided function for
1788 // the answer and post a send for it.
1789 auto &request_buffer = requests_buffers[index];
1790 request_buffer =
1791 (answer_request ?
1792 Utilities::pack(answer_request(other_rank,
1793 Utilities::unpack<RequestType>(
1794 buffer_recv, false)),
1795 false) :
1796 std::vector<char>());
1797
1798 ierr = MPI_Isend(request_buffer.data(),
1799 request_buffer.size(),
1800 MPI_CHAR,
1801 other_rank,
1802 tag_deliver,
1803 comm,
1804 &send_answer_requests[index]);
1805 AssertThrowMPI(ierr);
1806# else
1807 (void)answer_request;
1808 (void)comm;
1809 (void)index;
1810# endif
1811 }
1812
1813
1814
1815 template <typename RequestType, typename AnswerType>
1816 void
1818 const unsigned int n_targets,
1819 const std::function<void(const unsigned int, const AnswerType &)>
1820 &process_answer,
1821 const MPI_Comm comm)
1822 {
1823# ifdef DEAL_II_WITH_MPI
1824 const int tag_deliver = Utilities::MPI::internal::Tags::
1826
1827 // We know how many targets we have sent requests to. These
1828 // targets will all eventually send us their responses, but
1829 // we need not process them in order -- rather, just see what
1830 // comes in and then look at message originators' ranks and
1831 // message sizes
1832 for (unsigned int i = 0; i < n_targets; ++i)
1833 {
1834 MPI_Status status;
1835 {
1836 const int ierr =
1837 MPI_Probe(MPI_ANY_SOURCE, tag_deliver, comm, &status);
1838 AssertThrowMPI(ierr);
1839 }
1840
1841 const auto other_rank = status.MPI_SOURCE;
1842 int message_size;
1843 {
1844 const int ierr = MPI_Get_count(&status, MPI_CHAR, &message_size);
1845 AssertThrowMPI(ierr);
1846 }
1847 std::vector<char> recv_buffer(message_size);
1848
1849 // Now actually receive the answer. Because the MPI_Probe
1850 // above blocks until we have a message, we know that the
1851 // following MPI_Recv call will immediately succeed.
1852 {
1853 const int ierr = MPI_Recv(recv_buffer.data(),
1854 recv_buffer.size(),
1855 MPI_CHAR,
1856 other_rank,
1857 tag_deliver,
1858 comm,
1859 MPI_STATUS_IGNORE);
1860 AssertThrowMPI(ierr);
1861 }
1862
1863 if (process_answer)
1864 process_answer(other_rank,
1865 Utilities::unpack<AnswerType>(recv_buffer, false));
1866 }
1867# else
1868 (void)n_targets;
1869 (void)process_answer;
1870 (void)comm;
1871# endif
1872 }
1873
1874
1875
1876 template <typename RequestType, typename AnswerType>
1877 void
1879 {
1880# ifdef DEAL_II_WITH_MPI
1881 // Finalize all MPI_Request objects for both the
1882 // send-request and receive-answer operations.
1883 if (send_request_requests.size() > 0)
1884 {
1885 const int ierr = MPI_Waitall(send_request_requests.size(),
1886 send_request_requests.data(),
1887 MPI_STATUSES_IGNORE);
1888 AssertThrowMPI(ierr);
1889 }
1890
1891 // Then also check the send-answer requests.
1892 if (send_answer_requests.size() > 0)
1893 {
1894 const int ierr = MPI_Waitall(send_answer_requests.size(),
1895 send_answer_requests.data(),
1896 MPI_STATUSES_IGNORE);
1897 AssertThrowMPI(ierr);
1898 }
1899# endif
1900 }
1901
1902
1903
1904 template <typename RequestType, typename AnswerType>
1905 std::vector<unsigned int>
1907 const std::vector<unsigned int> &targets,
1908 const std::function<RequestType(const unsigned int)> &create_request,
1909 const std::function<AnswerType(const unsigned int, const RequestType &)>
1910 &answer_request,
1911 const std::function<void(const unsigned int, const AnswerType &)>
1912 &process_answer,
1913 const MPI_Comm comm)
1914 {
1916 ExcMessage("You shouldn't use the 'Serial' class on "
1917 "communicators that have more than one process "
1918 "associated with it."));
1919
1920 // The only valid target for a serial program is itself.
1921 if (targets.size() != 0)
1922 {
1923 Assert(targets.size() == 1,
1924 ExcMessage(
1925 "On a single process, the only valid target "
1926 "is process zero (the process itself), which can only be "
1927 "listed once."));
1928 AssertDimension(targets[0], 0);
1929
1930 // Since the caller indicates that there is a target, and since we
1931 // know that it is the current process, let the process send
1932 // something to itself.
1933 const RequestType request =
1934 (create_request ? create_request(0) : RequestType());
1935 const AnswerType answer =
1936 (answer_request ? answer_request(0, request) : AnswerType());
1937
1938 if (process_answer)
1939 process_answer(0, answer);
1940 }
1941
1942 return targets; // nothing to do
1943 }
1944
1945
1946
1947 template <typename RequestType, typename AnswerType>
1948 std::vector<unsigned int>
1950 const std::vector<unsigned int> &targets,
1951 const std::function<RequestType(const unsigned int)> &create_request,
1952 const std::function<AnswerType(const unsigned int, const RequestType &)>
1953 &answer_request,
1954 const std::function<void(const unsigned int, const AnswerType &)>
1955 &process_answer,
1956 const MPI_Comm comm)
1957 {
1958 // Depending on the number of processes we switch between
1959 // implementations. We reduce the threshold for debug mode to be
1960 // able to test also the non-blocking implementation. This feature
1961 // is tested by:
1962 // tests/multigrid/transfer_matrix_free_06.with_mpi=true.with_p4est=true.with_trilinos=true.mpirun=10.output
1963
1964 const unsigned int n_procs = (Utilities::MPI::job_supports_mpi() ?
1966 1);
1967# ifdef DEAL_II_WITH_MPI
1968# ifdef DEBUG
1969 if (n_procs > 10)
1970# else
1971 if (n_procs > 99)
1972# endif
1973 consensus_algo.reset(new NBX<RequestType, AnswerType>());
1974 else
1975# endif
1976 if (n_procs > 1)
1977 consensus_algo.reset(new PEX<RequestType, AnswerType>());
1978 else
1979 consensus_algo.reset(new Serial<RequestType, AnswerType>());
1980
1981 return consensus_algo->run(
1982 targets, create_request, answer_request, process_answer, comm);
1983 }
1984
1985
1986 } // namespace ConsensusAlgorithms
1987 } // end of namespace MPI
1988} // end of namespace Utilities
1989
1990#endif // DOXYGEN
1991
1992
1994
1995#endif
const char * get_exc_name() const
virtual const char * what() const noexcept override
virtual std::vector< unsigned int > run(const std::vector< unsigned int > &targets, const std::function< RequestType(const unsigned int)> &create_request, const std::function< AnswerType(const unsigned int, const RequestType &)> &answer_request, const std::function< void(const unsigned int, const AnswerType &)> &process_answer, const MPI_Comm comm)=0
void start_communication(const std::vector< unsigned int > &targets, const std::function< RequestType(const unsigned int)> &create_request, const MPI_Comm comm)
void maybe_answer_one_request(const std::function< AnswerType(const unsigned int, const RequestType &)> &answer_request, const MPI_Comm comm)
std::vector< std::unique_ptr< std::vector< char > > > request_buffers
virtual std::vector< unsigned int > run(const std::vector< unsigned int > &targets, const std::function< RequestType(const unsigned int)> &create_request, const std::function< AnswerType(const unsigned int, const RequestType &)> &answer_request, const std::function< void(const unsigned int, const AnswerType &)> &process_answer, const MPI_Comm comm) override
void clean_up_and_end_communication(const MPI_Comm comm)
bool all_locally_originated_receives_are_completed(const std::function< void(const unsigned int, const AnswerType &)> &process_answer, const MPI_Comm comm)
std::vector< std::unique_ptr< MPI_Request > > request_requests
std::vector< std::vector< char > > send_buffers
void signal_finish(const MPI_Comm comm)
unsigned int start_communication(const std::vector< unsigned int > &targets, const std::function< RequestType(const unsigned int)> &create_request, const MPI_Comm comm)
std::vector< std::vector< char > > requests_buffers
virtual std::vector< unsigned int > run(const std::vector< unsigned int > &targets, const std::function< RequestType(const unsigned int)> &create_request, const std::function< AnswerType(const unsigned int, const RequestType &)> &answer_request, const std::function< void(const unsigned int, const AnswerType &)> &process_answer, const MPI_Comm comm) override
std::vector< std::vector< char > > send_buffers
void answer_one_request(const unsigned int index, const std::function< AnswerType(const unsigned int, const RequestType &)> &answer_request, const MPI_Comm comm)
std::vector< std::vector< char > > recv_buffers
void process_incoming_answers(const unsigned int n_targets, const std::function< void(const unsigned int, const AnswerType &)> &process_answer, const MPI_Comm comm)
virtual std::vector< unsigned int > run(const std::vector< unsigned int > &targets, const std::function< RequestType(const unsigned int)> &create_request, const std::function< AnswerType(const unsigned int, const RequestType &)> &answer_request, const std::function< void(const unsigned int, const AnswerType &)> &process_answer, const MPI_Comm comm) override
std::shared_ptr< Interface< RequestType, AnswerType > > consensus_algo
virtual std::vector< unsigned int > run(const std::vector< unsigned int > &targets, const std::function< RequestType(const unsigned int)> &create_request, const std::function< AnswerType(const unsigned int, const RequestType &)> &answer_request, const std::function< void(const unsigned int, const AnswerType &)> &process_answer, const MPI_Comm comm) override
#define DEAL_II_NAMESPACE_OPEN
Definition config.h:38
constexpr bool running_in_debug_mode()
Definition config.h:76
#define DEAL_II_NAMESPACE_CLOSE
Definition config.h:39
#define Assert(cond, exc)
#define AssertDimension(dim1, dim2)
#define AssertThrowMPI(error_code)
#define AssertIndexRange(index, range)
static ::ExceptionBase & ExcMessage(std::string arg1)
const MPI_Comm comm
Definition mpi.cc:912
const unsigned int n_procs
Definition mpi.cc:923
std::vector< unsigned int > selector(const std::vector< unsigned int > &targets, const std::function< RequestType(const unsigned int)> &create_request, const std::function< AnswerType(const unsigned int, const RequestType &)> &answer_request, const std::function< void(const unsigned int, const AnswerType &)> &process_answer, const MPI_Comm comm)
std::vector< unsigned int > serial(const std::vector< unsigned int > &targets, const std::function< RequestType(const unsigned int)> &create_request, const std::function< AnswerType(const unsigned int, const RequestType &)> &answer_request, const std::function< void(const unsigned int, const AnswerType &)> &process_answer, const MPI_Comm comm)
std::vector< unsigned int > nbx(const std::vector< unsigned int > &targets, const std::function< RequestType(const unsigned int)> &create_request, const std::function< AnswerType(const unsigned int, const RequestType &)> &answer_request, const std::function< void(const unsigned int, const AnswerType &)> &process_answer, const MPI_Comm comm)
std::vector< unsigned int > pex(const std::vector< unsigned int > &targets, const std::function< RequestType(const unsigned int)> &create_request, const std::function< AnswerType(const unsigned int, const RequestType &)> &answer_request, const std::function< void(const unsigned int, const AnswerType &)> &process_answer, const MPI_Comm comm)
@ consensus_algorithm_nbx_process_deliver
ConsensusAlgorithms::NBX::process.
Definition mpi_tags.h:89
@ consensus_algorithm_pex_process_deliver
ConsensusAlgorithms::PEX::process.
Definition mpi_tags.h:94
@ consensus_algorithm_nbx_answer_request
ConsensusAlgorithms::NBX::process.
Definition mpi_tags.h:87
@ consensus_algorithm_pex_answer_request
ConsensusAlgorithms::PEX::process.
Definition mpi_tags.h:92
unsigned int n_mpi_processes(const MPI_Comm mpi_communicator)
Definition mpi.cc:103
unsigned int this_mpi_process(const MPI_Comm mpi_communicator)
Definition mpi.cc:118
bool job_supports_mpi()
Definition mpi.cc:680
unsigned int compute_n_point_to_point_communications(const MPI_Comm mpi_comm, const std::vector< unsigned int > &destinations)
Definition mpi.cc:371
std::size_t pack(const T &object, std::vector< char > &dest_buffer, const bool allow_compression=true)
Definition utilities.h:1352
STL namespace.