2#include <boost/container/flat_map.hpp>
3#include <boost/io/ios_state.hpp>
36 communication->send(
static_cast<int>(m.size()), rankReceiver);
38 for (
auto const &i : m) {
39 auto const &rank = i.first;
40 auto const &indices = i.second;
41 communication->send(rank, rankReceiver);
42 communication->sendRange(indices, rankReceiver);
52 communication->receive(size, rankSender);
56 communication->receive(rank, rankSender);
64 communication->broadcast(
static_cast<int>(m.size()));
66 for (
auto const &i : m) {
67 auto const &rank = i.first;
68 auto const &indices = i.second;
69 communication->broadcast(rank);
70 communication->broadcast(indices);
80 communication->broadcast(size, rankBroadcaster);
84 communication->broadcast(rank, rankBroadcaster);
85 communication->broadcast(m[rank], rankBroadcaster);
100void print(std::map<
int, std::vector<int>>
const &m)
102 std::ostringstream oss;
107 for (
auto &j : i.second) {
108 oss << i.first <<
":" << j <<
'\n';
124 std::cout << oss.str();
134 size_t maximum = std::numeric_limits<size_t>::min();
135 size_t minimum = std::numeric_limits<size_t>::max();
139 maximum = std::max(maximum,
static_cast<size_t>(size));
140 minimum = std::min(minimum,
static_cast<size_t>(size));
150 maximum = std::max(maximum,
static_cast<size_t>(size));
151 minimum = std::min(minimum,
static_cast<size_t>(size));
156 if (minimum > maximum)
159 auto average =
static_cast<double>(total);
164 boost::io::ios_all_saver ias{std::cout};
165 std::cout << std::fixed << std::setprecision(3)
166 <<
"Number of Communication Partners per Interface Process:"
168 <<
" Total: " << total <<
"\n"
169 <<
" Maximum: " << maximum <<
"\n"
170 <<
" Minimum: " << minimum <<
"\n"
171 <<
" Average: " << average <<
"\n"
172 <<
"Number of Interface Processes: " << count <<
"\n"
185 size += i.second.size();
190 size_t maximum = std::numeric_limits<size_t>::min();
191 size_t minimum = std::numeric_limits<size_t>::max();
195 maximum = std::max(maximum,
static_cast<size_t>(size));
196 minimum = std::min(minimum,
static_cast<size_t>(size));
207 maximum = std::max(maximum,
static_cast<size_t>(size));
208 minimum = std::min(minimum,
static_cast<size_t>(size));
214 if (minimum > maximum)
217 auto average =
static_cast<double>(total);
222 boost::io::ios_all_saver ias{std::cout};
223 std::cout << std::fixed << std::setprecision(3)
224 <<
"Number of LVDIs per Interface Process:"
226 <<
" Total: " << total <<
'\n'
227 <<
" Maximum: " << maximum <<
'\n'
228 <<
" Minimum: " << minimum <<
'\n'
229 <<
" Average: " << average <<
'\n'
230 <<
"Number of Interface Processes: " << count <<
'\n'
261std::map<int, std::vector<int>> buildCommunicationMap(
268 auto iterator = thisVertexDistribution.find(thisRank);
269 if (iterator == thisVertexDistribution.end()) {
273 std::map<int, std::vector<int>> communicationMap;
276 PRECICE_ASSERT(std::is_sorted(iterator->second.begin(), iterator->second.end()));
279 for (
const auto &[rank, vertices] : otherVertexDistribution) {
289 if (iterator->second.empty() || vertices.empty() || (vertices.back() < iterator->second.at(0)) || (vertices.at(0) > iterator->second.back())) {
294 std::vector<int> inters;
298 vertices.begin(), vertices.end(),
299 std::back_inserter(inters));
301 if (!inters.empty()) {
302 communicationMap.insert({rank, std::move(inters)});
305 return communicationMap;
312 _communicationFactory(std::move(communicationFactory))
316PointToPointCommunication::~PointToPointCommunication()
322bool PointToPointCommunication::isConnected()
const
327void PointToPointCommunication::acceptConnection(std::string
const &acceptorName,
328 std::string
const &requesterName)
337 PRECICE_DEBUG(
"Exchange vertex distribution between both primary ranks");
338 Event e0(
"m2n.exchangeVertexDistribution");
340 auto c = _communicationFactory->newCommunication();
353 _mesh->setVertexDistribution(vertexDistribution);
376 std::map<int, std::vector<int>> communicationMap = m2n::buildCommunicationMap(
377 vertexDistribution, requesterVertexDistribution);
383 print(communicationMap);
387#ifdef P2P_LCM_PRINT_STATS
393 Event e4(
"m2n.createCommunications");
394 e4.addData(
"Connections", communicationMap.size());
395 if (communicationMap.empty()) {
401 _communication = _communicationFactory->newCommunication();
406 _communication->acceptConnectionAsServer(acceptorName,
410 communicationMap.size());
413 for (
auto const &comMap : communicationMap) {
414 int globalRequesterRank = comMap.first;
415 auto indices = std::move(communicationMap[globalRequesterRank]);
417 _mappings.push_back({globalRequesterRank, std::move(indices),
com::PtrRequest(), {}});
423void PointToPointCommunication::acceptPreConnection(std::string
const &acceptorName,
424 std::string
const &requesterName)
429 const std::vector<int> &localConnectedRanks = _mesh->getConnectedRanks();
431 if (localConnectedRanks.empty()) {
436 _communication = _communicationFactory->newCommunication();
438 _communication->acceptConnectionAsServer(
443 localConnectedRanks.size());
445 _connectionDataVector.reserve(localConnectedRanks.size());
447 for (
int connectedRank : localConnectedRanks) {
454void PointToPointCommunication::requestConnection(std::string
const &acceptorName,
455 std::string
const &requesterName)
464 PRECICE_DEBUG(
"Exchange vertex distribution between both primary ranks");
465 Event e0(
"m2n.exchangeVertexDistribution");
467 auto c = _communicationFactory->newCommunication();
468 c->requestConnection(acceptorName, requesterName,
469 "TMP-PRIMARYCOM-" + _mesh->getName(),
481 _mesh->setVertexDistribution(vertexDistribution);
504 std::map<int, std::vector<int>> communicationMap = m2n::buildCommunicationMap(
505 vertexDistribution, acceptorVertexDistribution);
511 print(communicationMap);
515#ifdef P2P_LCM_PRINT_STATS
521 Event e4(
"m2n.createCommunications");
522 e4.addData(
"Connections", communicationMap.size());
523 if (communicationMap.empty()) {
528 std::vector<com::PtrRequest> requests;
529 requests.reserve(communicationMap.size());
531 std::set<int> acceptingRanks;
532 for (
auto &i : communicationMap)
533 acceptingRanks.emplace(i.first);
536 _communication = _communicationFactory->newCommunication();
541 _communication->requestConnectionAsClient(acceptorName, requesterName,
546 for (
auto &i : communicationMap) {
547 auto globalAcceptorRank = i.first;
548 auto indices = std::move(i.second);
550 _mappings.push_back({globalAcceptorRank, std::move(indices),
com::PtrRequest(), {}});
556void PointToPointCommunication::requestPreConnection(std::string
const &acceptorName,
557 std::string
const &requesterName)
562 std::vector<int> localConnectedRanks = _mesh->getConnectedRanks();
564 if (localConnectedRanks.empty()) {
569 std::vector<com::PtrRequest> requests;
570 requests.reserve(localConnectedRanks.size());
571 _connectionDataVector.reserve(localConnectedRanks.size());
573 std::set<int> acceptingRanks(localConnectedRanks.begin(), localConnectedRanks.end());
575 _communication = _communicationFactory->newCommunication();
576 _communication->requestConnectionAsClient(acceptorName, requesterName,
580 for (
auto &connectedRank : localConnectedRanks) {
586void PointToPointCommunication::completeSecondaryRanksConnection()
590 for (
auto &i : _connectionDataVector) {
591 _mappings.push_back({i.remoteRank, std::move(localCommunicationMap[i.remoteRank]), i.request, {}});
595void PointToPointCommunication::closeConnection()
599 if (not isConnected())
602 checkBufferedRequests(
true);
604 _communication.reset();
606 _connectionDataVector.clear();
610void PointToPointCommunication::send(precice::span<double const> itemsToSend,
int valueDimension)
613 if (_mappings.empty() || itemsToSend.
empty()) {
617 for (
auto &
mapping : _mappings) {
618 auto buffer = std::make_shared<std::vector<double>>();
619 buffer->reserve(
mapping.indices.size() * valueDimension);
620 for (
auto index :
mapping.indices) {
621 for (
int d = 0; d < valueDimension; ++d) {
622 buffer->push_back(itemsToSend[index * valueDimension + d]);
626 bufferedRequests.emplace_back(request, buffer);
628 checkBufferedRequests(
false);
631void PointToPointCommunication::receive(precice::span<double> itemsToReceive,
int valueDimension)
633 if (_mappings.empty() || itemsToReceive.
empty()) {
637 std::fill(itemsToReceive.
begin(), itemsToReceive.
end(), 0.0);
639 for (
auto &
mapping : _mappings) {
640 mapping.recvBuffer.resize(
mapping.indices.size() * valueDimension);
644 for (
auto &
mapping : _mappings) {
648 for (
auto index :
mapping.indices) {
649 for (
int d = 0; d < valueDimension; ++d) {
650 itemsToReceive[index * valueDimension + d] +=
mapping.recvBuffer[i * valueDimension + d];
657void PointToPointCommunication::broadcastSend(
int itemToSend)
659 for (
auto &connectionData : _connectionDataVector) {
660 _communication->send(itemToSend, connectionData.remoteRank);
664void PointToPointCommunication::broadcastReceiveAll(std::vector<int> &itemToReceive)
668 for (
auto &connectionData : _connectionDataVector) {
669 _communication->receive(data, connectionData.remoteRank);
670 itemToReceive.push_back(data);
674void PointToPointCommunication::broadcastSendMesh()
676 for (
auto &connectionData : _connectionDataVector) {
677 com::sendMesh(*_communication, connectionData.remoteRank, *_mesh);
681void PointToPointCommunication::broadcastReceiveAllMesh()
683 for (
auto &connectionData : _connectionDataVector) {
688void PointToPointCommunication::scatterAllCommunicationMap(CommunicationMap &localCommunicationMap)
690 for (
auto &connectionData : _connectionDataVector) {
691 _communication->sendRange(localCommunicationMap[connectionData.remoteRank], connectionData.remoteRank);
695void PointToPointCommunication::gatherAllCommunicationMap(CommunicationMap &localCommunicationMap)
697 for (
auto &connectionData : _connectionDataVector) {
698 localCommunicationMap[connectionData.remoteRank] = _communication->receiveRange(connectionData.remoteRank,
com::asVector<int>);
702void PointToPointCommunication::checkBufferedRequests(
bool blocking)
706 for (
auto it = bufferedRequests.begin(); it != bufferedRequests.end();) {
707 if (it->first->test())
708 it = bufferedRequests.erase(it);
712 if (bufferedRequests.empty())
715 std::this_thread::yield();
#define PRECICE_DEBUG(...)
#define PRECICE_TRACE(...)
#define PRECICE_ASSERT(...)
Interface for all distributed solver to solver communication classes.
PointToPointCommunication(com::PtrCommunicationFactory communicationFactory, mesh::PtrMesh mesh)
std::map< Rank, std::vector< VertexID > > CommunicationMap
A mapping from remote local ranks to the IDs that must be communicated.
std::map< Rank, std::vector< VertexID > > VertexDistribution
A mapping from rank to used (not necessarily owned) vertex IDs.
A C++ 11 implementation of the non-owning C++20 std::span type.
PRECICE_SPAN_NODISCARD constexpr bool empty() const noexcept
constexpr iterator begin() const noexcept
constexpr iterator end() const noexcept
static Rank getRank()
Current rank.
static bool isPrimary()
True if this process is running the primary rank.
static auto allSecondaryRanks()
Returns an iterable range over salve ranks [1, _size).
static bool isSecondary()
True if this process is running a secondary rank.
static com::PtrCommunication & getCommunication()
Intra-participant communication.
void sendMesh(Communication &communication, int rankReceiver, const mesh::Mesh &mesh)
constexpr auto asVector
Allows to use Communication::AsVectorTag in a less verbose way.
std::shared_ptr< Request > PtrRequest
void receiveMesh(Communication &communication, int rankSender, mesh::Mesh &mesh)
virtual void closeConnection()=0
Disconnects from communication space, i.e. participant.
std::shared_ptr< CommunicationFactory > PtrCommunicationFactory
std::shared_ptr< Communication > PtrCommunication
constexpr auto data(C &c) -> decltype(c.data())
contains the logic of the parallel communication between participants.
void print(std::map< int, std::vector< int > > const &m)
void broadcastReceive(mesh::Mesh::VertexDistribution &m, int rankBroadcaster, const com::PtrCommunication &communication=utils::IntraComm::getCommunication())
void printCommunicationPartnerCountStats(std::map< int, std::vector< int > > const &m)
void printLocalIndexCountStats(std::map< int, std::vector< int > > const &m)
void receive(mesh::Mesh::VertexDistribution &m, int rankSender, const com::PtrCommunication &communication)
void broadcast(mesh::Mesh::VertexDistribution &m)
void send(mesh::Mesh::VertexDistribution const &m, int rankReceiver, const com::PtrCommunication &communication)
void broadcastSend(mesh::Mesh::VertexDistribution const &m, const com::PtrCommunication &communication=utils::IntraComm::getCommunication())
contains data mapping from points to meshes.
std::shared_ptr< Mesh > PtrMesh
static constexpr SynchronizeTag Synchronize
Convenience instance of the SynchronizeTag.
void set_intersection_indices(InputIt1 ref1, InputIt1 first1, InputIt1 last1, InputIt2 first2, InputIt2 last2, OutputIt d_first)
This function is by and large the same as std::set_intersection(). The only difference is that we don...