OPAL (Object Oriented Parallel Accelerator Library) 2024.2
OPAL
AsyncSendBuffers.h
Go to the documentation of this file.
1//
2// Class AsyncSendBuffer and AsyncSendBuffers
3//
4// Copyright (c) 2010 - 2013, Yves Ineichen, ETH Zürich
5// All rights reserved
6//
7// Implemented as part of the PhD thesis
8// "Toward massively parallel multi-objective optimization with application to
9// particle accelerators" (https://doi.org/10.3929/ethz-a-009792359)
10//
11// This file is part of OPAL.
12//
13// OPAL is free software: you can redistribute it and/or modify
14// it under the terms of the GNU General Public License as published by
15// the Free Software Foundation, either version 3 of the License, or
16// (at your option) any later version.
17//
18// You should have received a copy of the GNU General Public License
19// along with OPAL. If not, see <https://www.gnu.org/licenses/>.
20//
21#include <algorithm>
22#include <memory>
23#include <sstream>
24#include <string>
25#include <vector>
26
27#include "mpi.h"
28
30
31public:
32 AsyncSendBuffer(std::ostringstream& os) {
33 this->size_req = new MPI_Request();
34 this->buffer_req = new MPI_Request();
35 this->buf_size = os.str().length();
36 buffer = new char[buf_size];
37 std::memcpy(buffer, os.str().c_str(), buf_size);
38 }
39
41 delete size_req;
42 delete buffer_req;
43 delete[] buffer;
44 }
45
46 bool hasCompleted() {
47 int bufferflag = 0;
48 MPI_Test(this->buffer_req, &bufferflag, MPI_STATUS_IGNORE);
49 if(bufferflag) {
50 int sizeflag = 0;
51 MPI_Test(this->buffer_req, &sizeflag, MPI_STATUS_IGNORE);
52 if(sizeflag) {
53 return true;
54 }
55 }
56 return false;
57 }
58
59 void send(int recv_rank, int size_tag, int data_tag, MPI_Comm comm) {
60 MPI_Isend(&buf_size, 1, MPI_LONG, recv_rank, size_tag, comm, size_req);
61 MPI_Isend(buffer, buf_size, MPI_CHAR, recv_rank, data_tag, comm, buffer_req);
62 }
63
64private:
65 // can't use smart pointers because MPI will hold last valid reference to
66 // pointer
67 MPI_Request *size_req;
68 MPI_Request *buffer_req;
69 char *buffer;
70
71 std::size_t buf_size;
72};
73
74
76
77public:
79
80 void insert(std::shared_ptr<AsyncSendBuffer> buf) {
81 collection_.push_back(buf);
82 }
83
84 void cleanup() {
85 collection_.erase(
86 std::remove_if(
87 collection_.begin(),
88 collection_.end(),
89 [](const std::shared_ptr<AsyncSendBuffer>& buf) {
90 return buf->hasCompleted();
91 }),
92 collection_.end());
93 }
94
95 std::size_t size() {
96 return collection_.size();
97 }
98
99private:
100 std::vector< std::shared_ptr<AsyncSendBuffer> > collection_;
101};
MPI_Request * buffer_req
MPI_Request * size_req
AsyncSendBuffer(std::ostringstream &os)
void send(int recv_rank, int size_tag, int data_tag, MPI_Comm comm)
std::size_t buf_size
void insert(std::shared_ptr< AsyncSendBuffer > buf)
std::size_t size()
std::vector< std::shared_ptr< AsyncSendBuffer > > collection_