MADNESS 0.10.1
cloud.h
Go to the documentation of this file.
1
2/**
3 \file cloud.h
4 \brief Declares the \c Cloud class for storing data and transfering them between worlds
5 \ingroup world
6
7*/
8
9/**
10 * TODO: - delete container record upon caching if container is replicated
11 */
12
13#ifndef SRC_MADNESS_WORLD_CLOUD_H_
14#define SRC_MADNESS_WORLD_CLOUD_H_
15
16
18#include<algorithm>
19#include<any>
20#include<atomic>
21#include<iomanip>
22#include<limits>
23#include<list>
24#include<mutex>
25
26
27/*!
28 \file cloud.h
29 \brief Defines and implements most of madness cloud storage
30
31 TODO: check use of preprocessor directives
32 TODO: clear cache in destructor won't work because no subworld is present -> must be explicitly called, error prone/
33*/
34
35namespace madness {
36
37 /// \brief A utility to get the name of a type as a string from chatGPT
38 template<typename T>
39 struct type_name {
40 static const char* value() { return typeid(T).name();}
41 };
42
43 template<>
44 struct type_name<Function<double,1>> { static const char* value() { return "Function<double,1>"; } };
45 template<>
46 struct type_name<Function<double,2>> { static const char* value() { return "Function<double,2>"; } };
47 template<>
48 struct type_name<Function<double,3>> { static const char* value() { return "Function<double,3>"; } };
49 template<>
50 struct type_name<Function<double,4>> { static const char* value() { return "Function<double,4>"; } };
51 template<>
52 struct type_name<Function<double,5>> { static const char* value() { return "Function<double,5>"; } };
53 template<>
54 struct type_name<Function<double,6>> { static const char* value() { return "Function<double,6>"; } };
55
56 template<>
57 struct type_name<std::vector<Function<double,1>>> { static const char* value() { return "std::vector<Function<double,1>>"; } };
58 template<>
59 struct type_name<std::vector<Function<double,2>>> { static const char* value() { return "std::vector<Function<double,2>>"; } };
60 template<>
61 struct type_name<std::vector<Function<double,3>>> { static const char* value() { return "std::vector<Function<double,3>>"; } };
62 template<>
63 struct type_name<std::vector<Function<double,4>>> { static const char* value() { return "std::vector<Function<double,4>>"; } };
64 template<>
65 struct type_name<std::vector<Function<double,5>>> { static const char* value() { return "std::vector<Function<double,5>>"; } };
66 template<>
67 struct type_name<std::vector<Function<double,6>>> { static const char* value() { return "std::vector<Function<double,6>>"; } };
68
69template<typename keyT>
70struct Recordlist {
71 std::list<keyT> list;
72
73 Recordlist() : list() {};
74
75 explicit Recordlist(const keyT &key) : list{key} {};
76
77 Recordlist(const Recordlist &other) : list(other.list) {};
78
80 for (auto &l2 : list2.list) list.push_back(l2);
81 return *this;
82 }
83
84 Recordlist &operator+=(const keyT &key) {
85 list.push_back(key);
86 return *this;
87 }
88
90 keyT key = list.front();
91 list.pop_front();
92 return key;
93 }
94
95 std::size_t size() const {
96 return list.size();
97 }
98
99 // if type provides id() member function (i.e. WorldObject) use that for hashing, otherwise use hash_value() for
100 // fundamental types (see worldhash.h)
101 template <typename T>
102 using member_id_t = decltype(std::declval<T>().id());
103
104 template <typename T>
106
107 // if type provides a hashing function use that, intrusive hashing, see worldhash.h
108 template <typename T>
109 using member_hash_t = decltype(std::declval<T>().hash());
110
111 template <typename T>
113
114 template<typename T, std::size_t NDIM>
115 static keyT compute_record(const Function<T,NDIM>& arg) {return hash_value(arg.get_impl()->id());}
116
117 template<typename T, std::size_t NDIM>
119
120 template<typename keyQ, typename valueT>
122
123 template<typename keyQ, typename valueT>
124 static keyT compute_record(const std::shared_ptr<WorldContainer<keyQ,valueT>>& arg) {return hash_value(arg->id());}
125
126 template<typename T, std::size_t NDIM>
127 static keyT compute_record(const std::shared_ptr<madness::FunctionImpl<T, NDIM>>& arg) {return hash_value(arg->id());}
128
129 template<typename T>
130 static keyT compute_record(const std::vector<T>& arg) {return hash_range(arg.begin(), arg.end());}
131
132 template<typename T>
133 static keyT compute_record(const Tensor<T>& arg) {return hash_value(arg.normf());}
134
135 template<typename T>
136 static keyT compute_record(const std::shared_ptr<T>& arg) {return compute_record(*arg);}
137
138 template<typename T>
139 static keyT compute_record(const T& arg) {
140 if constexpr (has_member_id<T>::value) {
141 return hash_value(arg.id());
142 } else if constexpr (std::is_pointer_v<T> && has_member_id<std::remove_pointer_t<T>>::value) {
143 return hash_value(arg->id());
144 } else {
145 // compute hash_code for fundamental types
146 std::size_t hashtype = typeid(T).hash_code();
147 hash_combine(hashtype,hash_value(arg));
148 return hashtype;
149 }
150 }
151
152
153 friend std::ostream &operator<<(std::ostream &os, const Recordlist &arg) {
154 using namespace madness::operators;
155 os << arg.list;
156 return os;
157 }
158
159};
160
161/// Process map for the cloud's batch container
162
163/// Routes a record to an explicitly assigned owner, falling back to a hash for any
164/// record that was never registered. `set_owner` is collective, so all ranks agree on
165/// routing without further communication. Unlike the default cloud container this map
166/// is never reset to local, which is what keeps owner pinning stable.
167template <typename keyT, typename hashfunT = Hash<keyT>>
169private:
170 const int nproc;
171 hashfunT hashfun;
172 std::shared_ptr<std::map<keyT, ProcessID>> table;
173
174public:
175 CloudOwnerPmap(World& world, const hashfunT& hf = hashfunT())
176 : nproc(world.mpi.nproc()), hashfun(hf), table(new std::map<keyT, ProcessID>()) {}
177
178 /// collective: every rank must register the same (key, owner) pair
179 void set_owner(const keyT& key, const ProcessID owner) { (*table)[key] = owner; }
180
181 ProcessID owner(const keyT& key) const override {
182 auto it = table->find(key);
183 if (it != table->end()) return it->second;
184 if (nproc == 1) return 0;
185 return hashfun(key) % nproc;
186 }
187
188 void print_table(const std::string& tag = "") const {
189 print("CloudOwnerPmap::table", tag, "size=", table->size(), "(nproc=", nproc, ")");
190 for (const auto& kv : *table) {
191 std::ostringstream os;
192 os << " key=0x" << std::hex << kv.first << std::dec << " owner=" << kv.second;
193 print(os.str());
194 }
195 }
196};
197
198/// BatchTransport holds a back-reference to its Cloud; its bodies are defined below the class
199class Cloud;
200
202/// the serialized batch travels by value through a Future, so it must be
203/// archive-serializable: a plain vector is, a shared_ptr is not
204using batch_bytesT = std::vector<unsigned char>;
205
206/// Point-to-point transfer of serialized function batches between universe ranks
207
208/// The bytes stream straight from the owner's local batch store to the requester by
209/// MPI point-to-point; they never ride inside an active-message payload, so there is
210/// no eager-buffer limit and no extra copy on the wire.
211///
212/// Both endpoints are posted from **comm-thread AM handlers** (`WorldObject::send` runs
213/// the member inline on the RMI receiver thread), never from worker tasks. That is what
214/// buys overlap under worker saturation: at a tight protocol every worker sits in the
215/// exchange kernel for a long time, so an endpoint posted as a *task* would queue behind
216/// it and the MPI op would not be posted until compute ended. On the comm thread the RMI
217/// loop's Testsome drives the rendezvous to completion during compute instead, leaving
218/// only the final await for the worker.
219///
220/// Wire protocol:
221/// 1. requester (worker): record a pending slot keyed by `tag`, send `on_trigger`
222/// 2. owner `on_trigger` (comm thread): Isend the local bytes, reply `on_reply` with the size
223/// 3. requester `on_reply` (comm thread): size the buffer, post the Irecv, enqueue `finish_recv`
224/// 4. requester `finish_recv` (worker): await the Irecv and set the future
225///
226/// The size travels in the reply of step 2 rather than a separate Isend so that step 3 can
227/// post the payload Irecv *during* compute; posting it at consume time would move the data
228/// transfer to post-compute and lose the overlap.
229class BatchTransport : public WorldObject<BatchTransport> {
230public:
231 /// tags live in the range MADNESS does not manage (safempi.h); 32767 is the
232 /// conservative MPI_TAG_UB floor
233 static constexpr int BATCH_TAG_BASE = 8192;
234 static constexpr int BATCH_TAG_CAP = 32767;
235
236 /// reply size meaning "the owner does not hold this record"; see on_trigger
237 static constexpr std::size_t BATCH_NOT_FOUND = ~std::size_t(0);
238
239 /// bytes per MPI message in a batch transfer; a larger payload is split into several
240
241 /// MPI byte counts are `int`, and batches do exceed 2 GiB: a 161-orbital run at k=10
242 /// measured 3.03 GiB, whose count narrowed to a negative int. MPI rejects that and
243 /// SafeMPI throws it on the comm thread, where nothing catches it.
244 static std::size_t batch_chunk_bytes() { return batch_chunk_bytes_; }
245
246 /// Set the chunk size, for tests only.
247
248 /// Collective and not thread-safe: call it on every rank before any transfer. It exists so
249 /// a unit test can reach the multi-chunk path with an affordable payload.
250 static void set_batch_chunk_bytes(const std::size_t n) {
251 MADNESS_CHECK_THROW(n > 0 and n <= std::size_t(std::numeric_limits<int>::max()),
252 "batch chunk size must be positive and fit an MPI count");
254 }
255
256 /// @param[in] universe the world the cloud lives in (collective construction)
257 /// @param[in] cloud back-reference used to read owner-local batch bytes
258 BatchTransport(World& universe, Cloud* cloud)
259 : WorldObject<BatchTransport>(universe), cloud_(cloud), next_tag_(0) {
260 this->process_pending();
261 }
262
263 /// Future to the serialized bytes of `record`, fetched from its owner
264
265 /// Resolves locally without MPI when this rank owns the record. The trigger is in
266 /// flight on return, so the round trip overlaps work until the future is consumed.
268
269private:
270 Cloud* cloud_; ///< back-reference (not owned)
271 std::atomic<int> next_tag_;
272
273 /// requester-side state for one outstanding receive. Held by shared_ptr so the buffer
274 /// address is stable for the Irecv and the entry can leave pending_ while finish_recv
275 /// still owns its reference.
276 struct PendingRecv {
278 int tag = -1;
279 bool not_found = false; ///< owner reported it does not hold the record
280 batch_bytesT buf; ///< sized in on_reply
281 std::vector<SafeMPI::Request> reqs; ///< one per chunk, posted in on_reply
282 Future<batch_bytesT> fut; ///< set by finish_recv
283 };
284 std::mutex pending_mtx_;
285 std::map<int, std::shared_ptr<PendingRecv>> pending_;
286
287 /// owner-side in-flight Isends, reaped lazily. Their buffers live in the cloud's
288 /// batch container and stay valid for the duration, so an un-reaped Isend is harmless.
289 std::mutex sends_mtx_;
290 std::list<SafeMPI::Request> sends_;
291
292 int alloc_tag() {
293 const int span = BATCH_TAG_CAP - BATCH_TAG_BASE + 1;
294 const int t = next_tag_.fetch_add(1);
295 return BATCH_TAG_BASE + ((t % span) + span) % span;
296 }
297
298 static inline std::size_t batch_chunk_bytes_ = std::size_t(1) << 30; ///< 1 GiB
299
300 void reap_sends();
301
302 /// owner side, comm thread: Isend the record's bytes, then reply with the count.
303 /// Must not throw: an exception here escapes every task-level handler.
304 void on_trigger(batch_keyT record, ProcessID requester, int tag);
305
306 /// requester side, comm thread: size the buffer, post the Irecvs, enqueue finish_recv
307
308 /// \param chunk the chunking the *owner* used, echoed back so the receiver never has to
309 /// infer it. See on_trigger.
310 void on_reply(int tag, std::size_t size, std::size_t chunk);
311
312 /// requester side, worker task: await the background-progressed Irecv and set the future
313 void finish_recv(int tag);
314};
315
316/// cloud class
317
318/// store and load data to/from the cloud into arbitrary worlds
319///
320/// Distributed data is always bound to a certain world. If it needs to be
321/// present in another world it can be serialized to the cloud and deserialized
322/// from there again. For an example see test_cloud.cc
323///
324/// Data is stored into a distributed container living in the universe.
325/// During storing a (replicated) list of records is returned that can be used to find the data
326/// in the container. If a combined object (a vector, tuple, etc) is stored a list of records
327/// will be generated. When loading the data from the world the record list will be used to
328/// deserialize all stored objects.
329///
330/// Note that there must be a fence after the destruction of subworld containers, as in:
331///
332/// create subworlds
333/// {
334/// dcT(subworld)
335/// do work
336/// }
337/// subworld.gop.fence();
338class Cloud {
339
340 bool debug = false; ///< prints debug output
341 bool is_replicated=false; ///< if contents of the container are replicated
342 bool is_rank_replicated=false; ///< if every rank holds its own copy of the container (see replicate())
343 bool dofence = true; ///< fences after load/store
344 bool force_load_from_cache = false; ///< forces load from cache (mainly for debugging)
345 bool use_cache=true;
346
347public:
348
349 typedef std::any cached_objT;
351 using valueT = std::vector<unsigned char>;
352 typedef std::map<keyT, cached_objT> cacheT;
354
356 StoreFunction, ///< store a madness function in the cloud -- can have a large memory impact
357 ///< equivalent to a deep copy
358 StoreFunctionPointer, ///< store the pointer to the function in the cloud.
359 ///< Return type still is a Function<T,NDIM> with a pointer to the universe function impl.
360 ///< equivalent to a shallow copy
361 };
362
363
364 friend std::ostream& operator<<(std::ostream& os, const StoragePolicy& sp) {
365 switch(sp) {
366 case StoreFunction: os << "Function"; break;
367 case StoreFunctionPointer: os << "FunctionPointer"; break;
368 default: os << "UnknownStoragePolicy"; break;
369 }
370 return os;
371 }
372
373 friend std::string to_string(const StoragePolicy sp) {
374 std::ostringstream os;
375 os << sp;
376 return os.str();
377 }
378
379private:
380 /// are the functions (WorldObjects) stored in the cloud or only pointers to them
382
383 /// cloud is a container: replication policy for the cloud container: distributed, node-replicated, rank-replicated
385
387 /// dedicated container for owner-pinned function batches; see store_batch / fetch_batch_p2p.
388 /// Uses CloudOwnerPmap so each batch record lives on an explicitly chosen rank.
389 std::shared_ptr<CloudOwnerPmap<keyT>> batch_pmap;
391 /// constructed after batch_container so it is destroyed first, as WorldObject lifetimes require
392 std::unique_ptr<BatchTransport> batch_transport_;
394 recordlistT local_list_of_container_keys; // a world-local list of keys occupied in container
395
396public:
397 std::list<WorldObjectBase*> world_object_base_list; // list of world objects stored in the cloud
398
399 template <typename T>
400 using member_cloud_serialize_t = decltype(std::declval<T>().cloud_store(std::declval<World&>(), std::declval<Cloud&>()));
401
402 template <typename T>
404
405public:
406
407 /// @param[in] universe the universe world
408 Cloud(madness::World &universe) : container(universe),
409 batch_pmap(new CloudOwnerPmap<keyT>(universe)),
410 batch_container(universe, batch_pmap),
411 batch_transport_(new BatchTransport(universe, this)),
412 reading_time(0l), copy_time(0l), writing_time(0l),
413 cache_reads(0l), cache_stores(0l) {
414 }
415
417 if ((not cached_objects.empty()) or (not local_list_of_container_keys.list.empty())) {
418 print("\nCloud::~Cloud(): cached_objects not empty, size=", cached_objects.size());
419 print("You need to call clear_cache(subworld) before destroying the cloud");
420 print("\n------------------------------\n");
421 std::string msg="deferred destruction of cloud with non-empty cache";
422 std::cerr << msg << std::endl;
423 }
424 }
425
426 void set_debug(bool value) {
427 debug = value;
428 }
429
430 /// dump the record->owner table on rank 0; pair with the caller's own task-to-rank
431 /// print to check that task assignment and batch routing agree
432 void print_batch_owner_map(World& universe, const std::string& tag = "") const {
433 if (universe.rank() != 0) return;
434 batch_pmap->print_table(tag);
435 }
436
437 void set_fence(bool value) {
438 dofence = value;
439 }
440
441 void set_force_load_from_cache(bool value) {
442 force_load_from_cache = value;
443 }
444
445 /// is the cloud container replicated: per rank, per node, or distributed
448 use_cache=false;
449 if (value == RankReplicated) use_cache=true;
450 }
451
452 /// is the cloud container replicated: per rank, per node, or distributed
456
459 if (disttype!=cloud_replication_policy) {
460 std::cout << "Cloud::validate_distribution(): distribution type mismatch, container is " << disttype
461 << " but cloud_replication_policy is " << cloud_replication_policy << std::endl;
462 return false;
463 }
464 return true;
465 }
466
467
468 /// storing policy refers to storing functions or pointers to functions
470 storage_policy = value;
471 }
472
473 /// storing policy refers to storing functions or pointers to functions
477
478 void print_size(World& universe) {
479 nlohmann::json stats=gather_memory_statistics(universe);
480 double byte2gbyte=1.0/(1024*1024*1024);
481 double global_memsize=stats["memory_global_GB"].template get<double>();
482 double max_record_size=stats["max_record_size"].template get<double>();
483 double min_memsize=stats["memory_min_GB"].template get<double>();
484 double max_memsize=stats["memory_max_GB"].template get<double>();
485 double global_size=stats["container_size_global"].template get<double>();
486
487 if (universe.rank()==0) {
488 print("Cloud memory:");
489 print(" replicated:",is_replicated);
490 print("size of cloud (total)");
491 print(" number of records: ",global_size);
492 print(" memory in GBytes: ",global_memsize);
493 print("size of cloud (average per node)");
494 print(" number of records: ",double(global_size)/universe.size());
495 print(" memory in GBytes: ",global_memsize/universe.size());
496 print("min/max of node");
497 print(" memory in GBytes: ",min_memsize,max_memsize);
498 print(" max record size in GBytes:",max_record_size*byte2gbyte);
499
500 }
501 }
502
503 /// return a json object with the cloud settings and statistics
504 nlohmann::json get_statistics(World& world) const {
505 nlohmann::json j;
506 { // settings
507 j["storage_policy"]=to_string(storage_policy);
508 j["cloud_replication_policy"]=to_string(cloud_replication_policy);
509 j["is_replicated"]=is_replicated;
510 j["local_cached_objects_size"]=cached_objects.size();
511 }
512 // timings
513 j.update(gather_timings(world));
514 j.update(gather_memory_statistics(world));
515 return j;
516
517 }
518
519 /// get size of the cloud container
520 nlohmann::json gather_memory_statistics(World &universe) const {
521
522 std::size_t memsize=0;
523 std::size_t max_record_size=0;
524 for (auto& item : container) {
525 memsize+=item.second.size();
526 max_record_size=std::max(max_record_size,item.second.size());
527 }
528 // batch records are held separately, and for a caller that stores batches they are
529 // the larger half; leaving them out would understate the cloud's footprint
530 std::size_t batch_memsize=0;
531 for (auto& item : batch_container) {
532 batch_memsize+=item.second.size();
533 max_record_size=std::max(max_record_size,item.second.size());
534 }
535 memsize+=batch_memsize;
536 std::size_t global_memsize=memsize;
537 std::size_t max_memsize=memsize;
538 std::size_t min_memsize=memsize;
539 double rss=madness::get_rss_usage_in_GB();
540 double rss_av=rss;
541 universe.gop.sum(global_memsize);
542 universe.gop.max(max_memsize);
543 universe.gop.max(max_record_size);
544 universe.gop.min(min_memsize);
545 universe.gop.max(rss);
546 universe.gop.sum(rss_av);
547 double byte2gbyte=1.0/(1024*1024*1024);
548
549 // convert type(container item).second to GB, i.e. number of bytes in the container to GB
550 double uchar2gbyte=byte2gbyte*sizeof(unsigned char);
551
552
553 auto local_size=container.size();
554 auto global_size=local_size;
555 universe.gop.sum(global_size);
556 auto batch_global_size=batch_container.size();
557 universe.gop.sum(batch_global_size);
558 std::size_t batch_global_memsize=batch_memsize;
559 universe.gop.sum(batch_global_memsize);
560 nlohmann::json j;
561 j["container_size_global"] = global_size;
562 j["batch_container_size_global"] = batch_global_size;
563 j["batch_memory_global_GB"] = batch_global_memsize*uchar2gbyte;
564 j["memory_global_GB"] = global_memsize*uchar2gbyte;
565 j["memory_min_GB"] = min_memsize*uchar2gbyte;
566 j["memory_max_GB"] = max_memsize*uchar2gbyte;
567 j["memory_rss_GB_max"] = rss;
568 j["memory_rss_GB_av"] = rss_av/universe.size();
569 j["max_record_size"] = max_record_size;
570 return j;
571 }
572
573 nlohmann::json gather_timings(World &universe) const {
574 double rtime_max = double(reading_time)*1.e-6;
575 double rtime_acc = double(reading_time)*1.e-6;
576 double rtime_av = double(reading_time)*1.e-6;
577 double ctime_max = double(copy_time)*1.e-6;
578 double ctime_acc = double(copy_time)*1.e-6;
579 double ctime_av = double(copy_time)*1.e-6;
580 double wtime = double(writing_time)*1.e-6;
581 double ptime = double(replication_time)*1.e-6;
582 double tptime = double(target_replication_time)*1.e-6;
583 universe.gop.max(rtime_max);
584 universe.gop.sum(rtime_acc);
585 rtime_av = rtime_acc/universe.size();
586 universe.gop.max(ctime_max);
587 universe.gop.sum(ctime_acc);
588 ctime_av = ctime_acc/universe.size();
589 universe.gop.max(wtime);
590 universe.gop.max(ptime);
591 universe.gop.max(tptime);
592 long creads = long(cache_reads);
593 long cstores = long(cache_stores);
594 universe.gop.sum(creads);
595 universe.gop.sum(cstores);
596 nlohmann::json j;
597 j["reading_time_max_s"] = rtime_max;
598 j["reading_time_acc_s"] = rtime_acc;
599 j["reading_time_av_s"] = rtime_av;
600 j["copy_time_max_s"] = ctime_max;
601 j["copy_time_acc_s"] = ctime_acc;
602 j["copy_time_av_s"] = ctime_av;
603 j["writing_time_s"] = wtime;
604 j["replication_time_s"] = ptime;
605 j["target_replication_time_s"] = tptime;
606 j["cache_reads"] = creads;
607 j["cache_stores"] = cstores;
608 return j;
609 }
610
611 /// backwards compatibility
612 void print_timings(World& universe) const {
613 print_timings(gather_timings(universe));
614 }
615
616 static void print_timings(const nlohmann::json timings) {
617 double rtime_max=timings["reading_time_max_s"].template get<double>();
618 double rtime_av=timings["reading_time_av_s"].template get<double>();
619 double rtime_acc=timings["reading_time_acc_s"].template get<double>();
620 // double ctime_max=timings["copy_time_max_s"].template get<double>();
621 // double ctime_av=timings["copy_time_av_s"].template get<double>();
622 // double ctime_acc=timings["copy_time_acc_s"].template get<double>();
623 double wtime=timings["writing_time_s"].template get<double>();
624 double ptime=timings["replication_time_s"].template get<double>();
625 double tptime=timings["target_replication_time_s"].template get<double>();
626 long creads=timings["cache_reads"].template get<long>();
627 long cstores=timings["cache_stores"].template get<long>();
628
629 auto precision = std::cout.precision();
630 std::cout << std::fixed << std::setprecision(1);
631 print("cloud storing wall time ", wtime);
632 print("cloud replication wall time ", ptime);
633 print("target replication wall time ", tptime);
634 print("cloud max reading time (all procs) ", rtime_max, std::defaultfloat);
635 print("cloud average reading cpu time (all procs) ", rtime_av, std::defaultfloat);
636 print("cloud accumulated reading cpu time (all procs) ", rtime_acc, std::defaultfloat);
637 std::cout << std::setprecision(precision) << std::scientific;
638 print("cloud cache stores ", long(cstores));
639 print("cloud cache loads ", long(creads));
640 }
641
642 static void print_memory_statistics(const nlohmann::json stats) {
643 double byte2gbyte=1.0/(1024*1024*1024);
644 double global_memsize=stats["memory_global_GB"].template get<double>();
645 double max_record_size=stats["max_record_size"].template get<double>();
646 double min_memsize=stats["memory_min_GB"].template get<double>();
647 double max_memsize=stats["memory_max_GB"].template get<double>();
648 double global_size=stats["container_size_global"].template get<double>();
649
650 print("Cloud memory:");
651 print(" size of cloud (total)");
652 print(" number of records: ",global_size);
653 print(" memory in GBytes: ",global_memsize);
654 // print(" size of cloud (average per node)");
655 // print(" number of records: ",double(global_size)/madness::world().size());
656 // print(" memory in GBytes: ",global_memsize*byte2gbyte/madness::world().size());
657 print(" min/max of node");
658 print(" memory in GBytes: ",min_memsize,max_memsize);
659 print(" max record size in GBytes:",max_record_size*byte2gbyte);
660 // the owner-pinned batches, reported separately because they are a different lifetime and
661 // usually the bulk of it. Zero unless some task stored batches.
662 const double b_size = stats.value("batch_container_size_global", 0.0);
663 if (b_size > 0.0) {
664 print(" owner-pinned batches");
665 print(" number of records: ", b_size);
666 print(" memory in GBytes: ", stats.value("batch_memory_global_GB", 0.0));
667 }
668 }
669
670 void clear_cache(World &subworld) {
671 cached_objects.clear();
672 local_list_of_container_keys.list.clear();
673 subworld.gop.fence();
674 }
675
676 void clear() {
677 container.clear();
678 // The owner-pinned batches too. They are not leaked without this -- an SCF run derives the
679 // same record keys each time, so the next application overwrites them -- but they are the
680 // largest thing the cloud holds, and without this they stay resident through everything
681 // that follows the exchange in an iteration, which is where the memory ceiling actually is.
682 batch_container.clear();
683 }
684
686 reading_time=0l;
687 copy_time=0l;
688 writing_time=0l;
689 writing_time1=0l;
692 cache_stores=0l;
693 cache_reads=0l;
694 }
695
696 /// functor to distribute/rank/node-replicate a function, passed in as a pointer to WorldObjectBase
697 template<typename T, std::size_t NDIM>
702 // figure out if wo is a FunctionImpl and do the distribution
703 if (auto fimpl=dynamic_cast<FunctionImpl<T, NDIM>*>(wo)) {
704 // fimpl->get_pmap()->print_data_sizes(world,"before distribution of function in cloud");
705 if (dt==RankReplicated) {
706 fimpl->replicate(false);
707 } else if (dt==NodeReplicated) {
708 // print("replicating function per node",fimpl);;
709 fimpl->replicate_on_hosts(true);
710 } else if (dt==Distributed) {
711 fimpl->undo_replicate(false);
712 } else {
713 MADNESS_EXCEPTION("unknown distribution type",1);
714 }
715 // fimpl->get_pmap()->print_data_sizes(world,"after distribution of function in cloud");
716 }
717 return 0;
718 }
719 };
720
721 /// distribute/node/rank replicate the targets of all world objects stored in the cloud
723 if (world_object_base_list.empty()) return;
724 World& world=world_object_base_list.front()->get_world();
725
726 for (auto wo : world_object_base_list) {
727 loop_types<DistributeFunctor, double, float, double_complex, float_complex>(std::tuple<DistributionType>(dt),wo);
728 // world.gop.fence();
729 }
730 world.gop.fence();
731
732 }
733
734 /// @param[in] world the subworld the objects are loaded to
735 /// @param[in] recordlist the list of records where the objects are stored
736
737 /// load a single object from the cloud, recordlist is kept unchanged
738 template<typename T>
739 T load(madness::World &world, const recordlistT recordlist) const {
740 recordlistT rlist = recordlist;
741 cloudtimer t(world, reading_time);
742
743 // forward_load will consume the recordlist while loading elements
744 return forward_load<T>(world, rlist);
745 }
746
747 /// similar to load, but will consume the recordlist
748
749 /// @param[in] world the subworld the objects are loaded to
750 /// @param[in] recordlist the list of records where the objects are stored
751 template<typename T>
752 T consuming_load(madness::World &world, recordlistT& recordlist) const {
753 cloudtimer t(world, reading_time);
754
755 // forward_load will consume the recordlist while loading elements
756 return forward_load<T>(world, recordlist);
757 }
758
759 /// load a single object from the cloud, recordlist is consumed while loading elements
760 template<typename T>
761 T forward_load(madness::World &world, recordlistT& recordlist) const {
762 // different objects are stored in different ways
763 // - tuples are split up into their components
764 // - classes with their own cloud serialization are stored using that
765 // - everything else is stored using their usual serialization
766 if constexpr (is_tuple<T>::value) {
767 return load_tuple<T>(world, recordlist);
768 } else if constexpr (has_cloud_serialize<T>::value) {
769 T target = allocator<T>(world);
770 target.cloud_load(world, *this, recordlist);
771 return target;
772 } else {
773 return do_load<T>(world, recordlist);
774 }
775 }
776
777 /// Register the owner of a batch record; local map insert, no communication
778
779 /// Collective in the same sense as store_batch: every rank must call it with an
780 /// identical (record, owner) pair or fetches will route inconsistently. Separating
781 /// registration from the payload lets all ranks replicate the routing while each
782 /// owner stores only its own bytes, over a size-1 subworld.
783 void register_batch_owner(const keyT record, const ProcessID owner) {
784 batch_pmap->set_owner(record, owner);
785 }
786
787 /// Store a batch of functions as one owner-pinned record
788
789 /// The whole vector -- its size and each function -- is serialized into a single
790 /// record in the batch container and routed to `owner`, so one batch is one record
791 /// with one owner. Must be called with identical (owner, record) on every rank of
792 /// `world`.
793 ///
794 /// @param[in] fence false lets a caller storing many batches emit one fence for all
795 /// of them; the collective gather inside the archive self-synchronizes
796 template<typename T, std::size_t NDIM>
797 keyT store_batch(madness::World& world, const std::vector<Function<T, NDIM>>& batch,
798 const ProcessID owner, const keyT record, const bool fence = true) {
799 if (is_replicated) {
800 print("Cloud contents are replicated and read-only!");
801 MADNESS_EXCEPTION("cloud error", 1);
802 }
803 batch_pmap->set_owner(record, owner);
804 cloudtimer t_batch(world, batch_store_time);
805 {
808 par.set_dofence(false);
809 std::size_t fsize = batch.size();
810 par & fsize;
811 for (std::size_t i = 0; i < fsize; ++i) par & batch[i];
812 }
813 if (dofence and fence) world.gop.fence();
814 return record;
815 }
816
817 /// @param[in] world presumably the universe
818 template<typename T>
820 if (is_replicated) {
821 print("Cloud contents are replicated and read-only!");
822 MADNESS_EXCEPTION("cloud error",1);
823 }
824 cloudtimer t(world,writing_time);
825
826 // different objects are stored in different ways
827 // - tuples are split up into their components
828 // - classes with their own cloud serialization are stored using that
829 // - everything else is stored using their usual serialization
830 recordlistT recordlist;
831 if constexpr (is_tuple<T>::value) {
832 recordlist+=store_tuple(world,source);
833 } else if constexpr (has_cloud_serialize<T>::value) {
834 recordlist+=source.cloud_store(world,*this);
835 } else {
836 recordlist+=store_other(world,source);
837 }
838 if (dofence) world.gop.fence();
839 return recordlist;
840 }
841
842 void replicate_according_to_policy(const std::size_t chunk_size=INT_MAX) {
844 // if (debug and (container.size() > 0)) print("no replication of container");
845 return;
846 }
848 replicate(chunk_size);
849 }
851 replicate_per_node(chunk_size);
852 }
853 else {
854 MADNESS_EXCEPTION("unknown replication policy",1);
855 }
856 container.get_world().gop.fence();
857 }
858
859 void replicate_per_node(const std::size_t chunk_size=INT_MAX) {
860 // this will fail if the container values are larger that 2GB
861 // need to reimplement that at some point
862 try {
863 double cpu0=cpu_time();
864 World& world=container.get_world();
865 world.gop.fence();
867 MADNESS_CHECK_THROW(not is_replicated,"cloud::replicate_per_node: container is already replicated");
868 container.replicate_on_hosts(true);
869 is_replicated=true;
870 world.gop.fence();
871 double cpu1=cpu_time();
872 if (debug and (world.rank()==0)) print("replication_per_node ended after ",cpu1-cpu0," seconds");
873 } catch (...) {
874 MADNESS_EXCEPTION("cloud replication_per_node failed, presumably because some data is larger than 2GB",1);
875 }
876 }
877
878 // replicates the contents of the container
879 void replicate(const std::size_t chunk_size=INT_MAX) {
880 MADNESS_CHECK_THROW(not is_replicated,"cloud::replicate_per_node: container is already replicated");
881
882 double cpu0=cpu_time();
883 World& world=container.get_world();
884 world.gop.fence();
886 container.reset_pmap_to_local();
887 is_replicated=true;
889
890 std::list<keyT> keylist;
891 for (auto it=container.begin(); it!=container.end(); ++it) {
892 keylist.push_back(it->first);
893 }
894
895 for (ProcessID rank=0; rank<world.size(); rank++) {
896 if (rank == world.rank()) {
897 std::size_t keylistsize = keylist.size();
898 world.mpi.Bcast(&keylistsize,sizeof(keylistsize),MPI_BYTE,rank);
899
900 for (auto key : keylist) {
902 bool found=container.find(acc,key);
903 MADNESS_CHECK(found);
904 auto data = acc->second;
905 std::size_t sz=data.size();
906
907 world.mpi.Bcast(&key,sizeof(key),MPI_BYTE,rank);
908 world.mpi.Bcast(&sz,sizeof(sz),MPI_BYTE,rank);
909
910 // if data is too large for MPI_INT break it into pieces to avoid integer overflow
911 for (std::size_t start=0; start<sz; start+=chunk_size) {
912 std::size_t remainder = std::min(sz - start, chunk_size);
913 world.mpi.Bcast(&data[start], remainder, MPI_BYTE, rank);
914 }
915
916 }
917 }
918 else {
919 std::size_t keylistsize = 0;
920 world.mpi.Bcast(&keylistsize,sizeof(keylistsize),MPI_BYTE,rank);
921 for (size_t i=0; i<keylistsize; i++) {
922 keyT key;
923 world.mpi.Bcast(&key,sizeof(key),MPI_BYTE,rank);
924 std::size_t sz = 0;
925 world.mpi.Bcast(&sz,sizeof(sz),MPI_BYTE,rank);
926 valueT data(sz);
927// world.mpi.Bcast(&data[0],sz,MPI_BYTE,rank);
928 for (std::size_t start=0; start<sz; start+=chunk_size) {
929 std::size_t remainder=std::min(sz-start,chunk_size);
930 world.mpi.Bcast(&data[start],remainder,MPI_BYTE,rank);
931 }
932
933 container.replace(key,data);
934 }
935 }
936 }
937 world.gop.fence();
938 double cpu1=cpu_time();
939 if (debug and (world.rank()==0)) print("replication ended after ",cpu1-cpu0," seconds");
940 }
941
942private:
943
944 mutable std::atomic<long> reading_time=0l; // in microseconds
945 mutable std::atomic<long> batch_store_time=0l; ///< store_batch wall time, microseconds
946 mutable std::atomic<long> batch_find_time=0l; ///< waiting on the p2p transfer, microseconds
947 mutable std::atomic<long> batch_deserialize_time=0l; ///< deserializing the bytes, microseconds
948public:
949 mutable std::atomic<long> copy_time=0l; // if pointers are stored in cloud, time to copy from universe to subworld
950 mutable std::atomic<long> target_replication_time=0l; // if pointers are stored in cloud, time to replicate targets
951private:
952 mutable std::atomic<long> writing_time=0l; // in microseconds
953 mutable std::atomic<long> writing_time1=0l; // in microseconds
954 mutable std::atomic<long> replication_time=0l; // in microseconds
955 mutable std::atomic<long> cache_reads=0l;
956 mutable std::atomic<long> cache_stores=0l;
957
958 template<typename> struct is_tuple : std::false_type { };
959 template<typename ...T> struct is_tuple<std::tuple<T...>> : std::true_type { };
960
961 template<typename Q> struct is_vector : std::false_type { };
962 template<typename Q> struct is_vector<std::vector<Q>> : std::true_type { };
963
964 template<typename T> using is_parallel_serializable_object = std::is_base_of<archive::ParallelSerializableObject,T>;
965
966 template<typename T> using is_world_constructible = std::is_constructible<T, World &>;
967
968public:
969 struct cloudtimer {
971 double wall0;
972 std::atomic<long> &rtime;
973
974 cloudtimer(World& world, std::atomic<long> &readtime) : world(world), wall0(wall_time()), rtime(readtime) {}
975
977 long deltatime=long((wall_time() - wall0) * 1000000l);
978 rtime += deltatime;
979 }
980 };
981private:
982
983 template<typename T>
984 void cache(madness::World &world, const T &obj, const keyT &record) const {
985 const_cast<cacheT &>(cached_objects).insert({record,std::make_any<T>(obj)});
986 }
987
988 /// load an object from the cache, record is unchanged
989 template<typename T>
990 T load_from_cache(madness::World &world, const keyT &record) const {
991 if (world.rank()==0) cache_reads++;
992 if (debug) print("loading", type_name<T>::value(), "from cache record", record, "to world", world.id());
993 if (auto obj = std::any_cast<T>(&cached_objects.find(record)->second)) return *obj;
994 MADNESS_EXCEPTION("failed to load from cloud-cache", 1);
995 T target = allocator<T>(world);
996 return target;
997 }
998
999 bool is_cached(const keyT &key) const {
1000 return (cached_objects.count(key) == 1);
1001 }
1002
1003public:
1004
1005 /// the owner of a batch record; a pmap lookup, no communication
1006 ProcessID batch_owner(const keyT record) const {
1007 return batch_pmap->owner(record);
1008 }
1009
1010private:
1011
1012 /// only the transport reads owner-local bytes; callers go through fetch_batch_p2p
1013 friend class BatchTransport;
1014
1015 /// bytes of a batch record held by this rank, empty if it holds none
1016
1017 /// Never throws: the callers are BatchTransport's comm-thread handlers, where an
1018 /// escaping exception bypasses every task-level handler and surfaces as an
1019 /// unattributable abort. A miss is reported to the requester instead, and raised
1020 /// there in task context.
1021 ///
1022 /// Emptiness is a sound "not found" marker because a stored batch always begins with
1023 /// its serialized element count, so a present record is never zero bytes.
1026 if (batch_container.find(acc, record)) return acc->second;
1027 return valueT();
1028 }
1029
1030 /// stable pointer and size of a local batch record, {nullptr,0} if this rank holds none
1031
1032 /// Lets the comm thread Isend without copying the payload. The accessor lock is
1033 /// released on return, but the address stays valid because batch records are neither
1034 /// erased nor mutated between the store and the end of the consuming operation.
1035 /// Never throws, for the reason given on try_get_local_batch_bytes.
1036 std::pair<const unsigned char*, std::size_t> try_get_local_batch_ptr(const keyT record) const {
1038 if (batch_container.find(acc, record))
1039 return {acc->second.data(), acc->second.size()};
1040 return {nullptr, 0};
1041 }
1042
1043public:
1044
1045 /// start fetching `record` from its owner; the trigger is in flight on return
1047 return batch_transport_->request(record);
1048 }
1049
1050 /// turn the bytes of a p2p transfer into the batch of functions
1051
1052 /// Blocks on `fut` only if the transfer has not landed yet. Runs in a task, which is
1053 /// where a missing record is reported so the failure is attributable.
1054 /// @param[in] cache_result default false: the cloud-side cache is not safe to keep
1055 /// across changes of the calling world, so opting in is the
1056 /// caller's decision
1057 template<typename T, std::size_t NDIM>
1058 std::vector<Function<T, NDIM>> deserialize_batch_p2p(madness::World& subworld,
1059 Future<batch_bytesT> fut, const keyT record,
1060 const bool cache_result = false) const {
1061 typedef std::vector<Function<T, NDIM>> vecfuncT;
1062 if (is_cached(record)) return load_from_cache<vecfuncT>(subworld, record);
1063 cloudtimer t(subworld, reading_time);
1064 const double w0 = wall_time();
1065 batch_bytesT bytes = fut.get();
1066 const double w1 = wall_time();
1067 batch_find_time += long((w1 - w0) * 1.e6);
1068 MADNESS_CHECK_THROW(not bytes.empty(),
1069 "deserialize_batch_p2p: the owner does not hold this batch record");
1070 vecfuncT batch;
1071 {
1074 std::size_t fsize = 0;
1075 par & fsize;
1076 batch.resize(fsize);
1077 for (std::size_t i = 0; i < fsize; ++i) par & batch[i];
1078 }
1079 batch_deserialize_time += long((wall_time() - w1) * 1.e6);
1080 if (use_cache and cache_result) cache(subworld, batch, record);
1081 return batch;
1082 }
1083
1084 /// fetch a batch stored by store_batch; resolves without MPI when this rank owns it
1085 template<typename T, std::size_t NDIM>
1086 std::vector<Function<T, NDIM>> fetch_batch_p2p(madness::World& subworld,
1087 const keyT record, const bool cache_result = false) const {
1088 typedef std::vector<Function<T, NDIM>> vecfuncT;
1089 if (is_cached(record)) return load_from_cache<vecfuncT>(subworld, record);
1090 return deserialize_batch_p2p<T, NDIM>(subworld, request_batch_bytes_async(record),
1091 record, cache_result);
1092 }
1093
1094private:
1095
1096 /// checks if a (universe) container record is used
1097
1098 /// currently implemented with a local copy of the recordlist, might be
1099 /// reimplemented with container.find(), which would include blocking communication.
1100 bool is_in_container(const keyT &key) const {
1101 auto it = std::find(local_list_of_container_keys.list.begin(),
1102 local_list_of_container_keys.list.end(), key);
1103 return it!=local_list_of_container_keys.list.end();
1104 }
1105
1106 template<typename T>
1107 T allocator(World &world) const {
1108 if constexpr (is_world_constructible<T>::value) {
1109 return T(world);
1110 } else {
1111 return T();
1112 }
1113 }
1114
1115 template<typename T>
1118 bool is_already_present= is_in_container(record);
1119 if (debug and world.rank()==0) {
1120 if (is_already_present) std::cout << "skipping ";
1121 if constexpr (Recordlist<keyT>::has_member_id<T>::value) {
1122 std::cout << "storing world object of " << type_name<T>::value() << "id " << source.id()
1123 << " to record " << record << std::endl;
1124 }
1125 std::cout << "storing object of " << type_name<T>::value() << " to record " << record << std::endl;
1126 }
1127 if constexpr (is_madness_function<T>::value) {
1128 if (source.is_compressed() and T::dimT>3) print("WARNING: storing compressed hi-dim `function");
1129 }
1130
1131 // scope is important because of destruction ordering of world objects and fence
1132 if (is_already_present) {
1133 if (world.rank()==0) cache_stores++;
1134 } else {
1135 cloudtimer t(world,writing_time1);
1139 if constexpr (is_madness_function<T>::value) {
1140 // store the pointer to the function, not the function itself
1141 par & source.get_impl();
1142 // store the pointer to the WorldObject in a list for later reference (replication/redistribution)
1143 WorldObjectBase* wobj=source.get_impl().get();
1144 world_object_base_list.push_back(wobj);
1145 } else {
1146 // store everything else
1147 par & source;
1148 }
1149 } else {
1150 // store everything else
1151 par & source;
1152 }
1154 }
1155 if (dofence) world.gop.fence();
1156 return recordlistT{record};
1157 }
1158
1159public:
1160 /// load a vector from the cloud, pop records from recordlist
1161 ///
1162 /// @param[inout] world destination world
1163 /// @param[inout] recordlist list of records to load from (reduced by the first few elements)
1164 template<typename T>
1165 typename std::enable_if<is_vector<T>::value, T>::type
1166 do_load(World &world, recordlistT &recordlist) const {
1167 std::size_t sz = do_load<std::size_t>(world, recordlist);
1168 T target(sz);
1169 for (std::size_t i = 0; i < sz; ++i) {
1170 target[i] = do_load<typename T::value_type>(world, recordlist);
1171 }
1172 return target;
1173 }
1174
1175 /// load a single object from the cloud, pop record from recordlist
1176 ///
1177 /// @param[inout] world destination world
1178 /// @param[inout] recordlist list of records to load from (reduced by the first element)
1179 template<typename T>
1180 typename std::enable_if<!is_vector<T>::value, T>::type
1181 do_load(World &world, recordlistT &recordlist) const {
1182 keyT record = recordlist.pop_front_and_return();
1184
1185 if (is_cached(record)) return load_from_cache<T>(world, record);
1186 if (debug) print("loading", type_name<T>::value(), "from container record", record, "to world", world.id());
1187 T target = allocator<T>(world);
1190 if constexpr (is_madness_function<T>::value) {
1192 // load the pointer to the function, not the function itself
1193 // this is important for large functions, as they are not replicated
1194 // and only copied to subworlds when needed
1195 try {
1197 std::shared_ptr<implT> impl;
1198 par & impl;
1199 target.set_impl(impl); // target now points to a universe function impl
1200 } catch (...) {
1201 {
1202 io_redirect_cout redirect;
1203 print("failed to load function pointer from cloud, maybe the target is out of scope?");
1204 print("record:", record, "world:", world.id());
1205 print("function type:", type_name<T>::value());
1206 print("\n");
1207 }
1208 MADNESS_EXCEPTION("load/store error of pointers in cloud", 1);
1209 }
1210 } else {
1211 // load everything else
1212 par & target;
1213 }
1214 } else {
1215 // load everything else
1216 par & target;
1217 }
1218
1219 if (use_cache) {
1220 cache(world, target, record);
1221 // a rank-replicated record is this rank's own copy, so release it now that it is cached. Not so when
1222 // node-replicated: the host's lowest rank owns the only copy, and erase() would remove it for every
1223 // rank on the host, which then fail with "record not found"
1224 if (is_rank_replicated) container.erase(record);
1225 }
1226
1227 return target;
1228 }
1229
1230public:
1231
1232 // overloaded
1233 template<typename T>
1234 recordlistT store_other(madness::World& world, const std::vector<T>& source) {
1235 if (debug and world.rank()==0)
1236 std::cout << "storing vector of " << type_name<T>::value() << " of size " << source.size() << std::endl;
1237 recordlistT l = store_other(world, source.size());
1238 for (const auto& s : source) l += store_other(world, s);
1239 if (dofence) world.gop.fence();
1240 if (debug and world.rank()==0) std::cout << "done with vector storing; container size "
1241 << container.size() << std::endl;
1242 return l;
1243 }
1244
1245 /// store a tuple in multiple records
1246 template<typename... Ts>
1247 recordlistT store_tuple(World &world, const std::tuple<Ts...> &input) {
1248 recordlistT v;
1249 auto storeaway = [&](const auto &arg) {
1250 v += this->store(world, arg);
1251 };
1252 auto l = [&](Ts const &... arg) {
1253 ((storeaway(arg)), ...);
1254 };
1255 std::apply(l, input);
1256 return v;
1257 }
1258
1259 /// load a tuple from the cloud, pop records from recordlist
1260 ///
1261 /// @param[inout] world destination world
1262 /// @param[inout] recordlist list of records to load from (reduced by the first few elements)
1263 template<typename T>
1264 T load_tuple(madness::World &world, recordlistT &recordlist) const {
1265 if (debug) std::cout << "loading tuple of type " << typeid(T).name() << " to world " << world.id() << std::endl;
1266 T target;
1267 std::apply([&](auto &&... args) {
1268 ((args = forward_load<typename std::remove_reference<decltype(args)>::type>(world, recordlist)), ...);
1269 }, target);
1270 return target;
1271 }
1272};
1273
1274// ---- BatchTransport bodies; they need the complete Cloud ----
1275
1277 World& u = this->get_world();
1278 const ProcessID owner = cloud_->batch_owner(record);
1279 if (owner == u.rank()) {
1280 // local: no MPI. An absent record yields empty bytes, which the consumer
1281 // reports in task context, exactly as for the remote path.
1283 }
1284 const int tag = alloc_tag();
1285 auto p = std::make_shared<PendingRecv>();
1286 p->owner = owner;
1287 p->tag = tag;
1288 {
1289 std::lock_guard<std::mutex> g(pending_mtx_);
1290 pending_[tag] = p;
1291 }
1292 // send, not add: the trigger runs inline on the owner's comm thread, so the owner
1293 // posts its Isend without queueing behind its own saturated workers
1294 this->send(owner, &BatchTransport::on_trigger, record, u.rank(), tag);
1295 return p->fut;
1296}
1297
1299 std::lock_guard<std::mutex> g(sends_mtx_);
1300 for (auto it = sends_.begin(); it != sends_.end(); ) {
1301 if (it->Test()) it = sends_.erase(it);
1302 else ++it;
1303 }
1304}
1305
1306inline void BatchTransport::on_trigger(batch_keyT record, ProcessID requester, int tag) {
1307 World& u = this->get_world();
1308 reap_sends();
1309 auto ptr_size = cloud_->try_get_local_batch_ptr(record);
1310 if (ptr_size.first == nullptr) {
1311 // Not held here -- reply with the sentinel and post nothing. Throwing instead
1312 // would escape this comm-thread handler past every task-level handler and abort
1313 // the run with no attribution; the requester raises it in task context.
1314 this->send(requester, &BatchTransport::on_reply, tag, BATCH_NOT_FOUND, std::size_t(0));
1315 return;
1316 }
1317 const std::size_t n = ptr_size.second;
1318 // Sent along rather than read again by the requester: the size is a process-local static, so
1319 // the two ends can disagree, and a receiver that guessed would post Irecvs that do not match
1320 // the messages -- MPI_ERR_TRUNCATE on the comm thread, where nothing catches it.
1321 const std::size_t chunk = batch_chunk_bytes_;
1322 // MPI does not overtake between messages sharing source, tag and communicator, so posting
1323 // the chunks in ascending offset order matches the requester's Irecvs pairwise without a
1324 // per-chunk tag. That does rely on one transfer per tag, which alloc_tag's span makes true.
1325 {
1326 std::lock_guard<std::mutex> g(sends_mtx_);
1327 for (std::size_t off = 0; off < n; off += chunk) {
1328 const std::size_t len = std::min(chunk, n - off);
1329 sends_.push_back(u.mpi.Isend(ptr_size.first + off, int(len), MPI_BYTE, requester, tag));
1330 }
1331 }
1332 // the size rides in the reply so the requester can post its Irecvs now, during its
1333 // own compute, and let the rendezvous finish in the background
1334 this->send(requester, &BatchTransport::on_reply, tag, n, chunk);
1335}
1336
1337inline void BatchTransport::on_reply(int tag, std::size_t size, std::size_t chunk) {
1338 World& u = this->get_world();
1339 std::shared_ptr<PendingRecv> p;
1340 {
1341 std::lock_guard<std::mutex> g(pending_mtx_);
1342 auto it = pending_.find(tag);
1343 MADNESS_CHECK(it != pending_.end());
1344 p = it->second;
1345 }
1346 if (size == BATCH_NOT_FOUND) {
1347 // no Isend was posted, so post no Irecv; empty bytes mark the miss
1348 {
1349 std::lock_guard<std::mutex> g(pending_mtx_);
1350 pending_.erase(tag);
1351 }
1352 p->not_found = true;
1353 p->fut.set(batch_bytesT());
1354 return;
1355 }
1356 p->buf.resize(size);
1357 // the owner's framing, not ours; see on_trigger. Positive whenever the record was found, and
1358 // unguarded because a throw on this thread is what the framing exists to avoid
1359 for (std::size_t off = 0; off < size; off += chunk) {
1360 const std::size_t len = std::min(chunk, size - off);
1361 p->reqs.push_back(u.mpi.Irecv(p->buf.data() + off, int(len), MPI_BYTE, p->owner, tag));
1362 }
1363 u.taskq.add(this, &BatchTransport::finish_recv, tag);
1364}
1365
1366inline void BatchTransport::finish_recv(int tag) {
1367 std::shared_ptr<PendingRecv> p;
1368 {
1369 std::lock_guard<std::mutex> g(pending_mtx_);
1370 auto it = pending_.find(tag);
1371 MADNESS_CHECK(it != pending_.end());
1372 p = it->second;
1373 pending_.erase(it);
1374 }
1375 // the comm thread has been progressing these Irecvs all along, so the awaits are short
1376 for (auto& r : p->reqs) World::await(r, true);
1377 p->fut.set(std::move(p->buf));
1378}
1379
1380} /* namespace madness */
1381
1382#endif /* SRC_MADNESS_WORLD_CLOUD_H_ */
Point-to-point transfer of serialized function batches between universe ranks.
Definition cloud.h:229
static void set_batch_chunk_bytes(const std::size_t n)
Set the chunk size, for tests only.
Definition cloud.h:250
void on_trigger(batch_keyT record, ProcessID requester, int tag)
Definition cloud.h:1306
static std::size_t batch_chunk_bytes_
1 GiB
Definition cloud.h:298
BatchTransport(World &universe, Cloud *cloud)
Definition cloud.h:258
Cloud * cloud_
back-reference (not owned)
Definition cloud.h:270
std::map< int, std::shared_ptr< PendingRecv > > pending_
Definition cloud.h:285
std::mutex sends_mtx_
Definition cloud.h:289
static constexpr int BATCH_TAG_CAP
Definition cloud.h:234
int alloc_tag()
Definition cloud.h:292
void finish_recv(int tag)
requester side, worker task: await the background-progressed Irecv and set the future
Definition cloud.h:1366
static constexpr std::size_t BATCH_NOT_FOUND
reply size meaning "the owner does not hold this record"; see on_trigger
Definition cloud.h:237
std::mutex pending_mtx_
Definition cloud.h:284
void reap_sends()
Definition cloud.h:1298
std::atomic< int > next_tag_
Definition cloud.h:271
void on_reply(int tag, std::size_t size, std::size_t chunk)
requester side, comm thread: size the buffer, post the Irecvs, enqueue finish_recv
Definition cloud.h:1337
static std::size_t batch_chunk_bytes()
bytes per MPI message in a batch transfer; a larger payload is split into several
Definition cloud.h:244
Future< batch_bytesT > request(batch_keyT record)
Future to the serialized bytes of record, fetched from its owner.
Definition cloud.h:1276
std::list< SafeMPI::Request > sends_
Definition cloud.h:290
static constexpr int BATCH_TAG_BASE
Definition cloud.h:233
Process map for the cloud's batch container.
Definition cloud.h:168
CloudOwnerPmap(World &world, const hashfunT &hf=hashfunT())
Definition cloud.h:175
ProcessID owner(const keyT &key) const override
Maps key to processor.
Definition cloud.h:181
void print_table(const std::string &tag="") const
Definition cloud.h:188
void set_owner(const keyT &key, const ProcessID owner)
collective: every rank must register the same (key, owner) pair
Definition cloud.h:179
hashfunT hashfun
Definition cloud.h:171
std::shared_ptr< std::map< keyT, ProcessID > > table
Definition cloud.h:172
const int nproc
Definition cloud.h:170
cloud class
Definition cloud.h:338
bool is_cached(const keyT &key) const
Definition cloud.h:999
void print_batch_owner_map(World &universe, const std::string &tag="") const
Definition cloud.h:432
Future< batch_bytesT > request_batch_bytes_async(const keyT record) const
start fetching record from its owner; the trigger is in flight on return
Definition cloud.h:1046
bool use_cache
Definition cloud.h:345
void clear()
Definition cloud.h:676
void replicate_per_node(const std::size_t chunk_size=INT_MAX)
Definition cloud.h:859
bool is_in_container(const keyT &key) const
checks if a (universe) container record is used
Definition cloud.h:1100
bool force_load_from_cache
forces load from cache (mainly for debugging)
Definition cloud.h:344
bool debug
prints debug output
Definition cloud.h:340
std::atomic< long > writing_time1
Definition cloud.h:953
nlohmann::json get_statistics(World &world) const
return a json object with the cloud settings and statistics
Definition cloud.h:504
std::enable_if< is_vector< T >::value, T >::type do_load(World &world, recordlistT &recordlist) const
Definition cloud.h:1166
std::atomic< long > batch_find_time
waiting on the p2p transfer, microseconds
Definition cloud.h:946
recordlistT store_other(madness::World &world, const std::vector< T > &source)
Definition cloud.h:1234
std::any cached_objT
Definition cloud.h:349
recordlistT store(madness::World &world, const T &source)
Definition cloud.h:819
T load_tuple(madness::World &world, recordlistT &recordlist) const
Definition cloud.h:1264
std::is_base_of< archive::ParallelSerializableObject, T > is_parallel_serializable_object
Definition cloud.h:964
~Cloud()
Definition cloud.h:416
madness::archive::ContainerRecordOutputArchive::keyT keyT
Definition cloud.h:350
std::atomic< long > replication_time
Definition cloud.h:954
Recordlist< keyT > recordlistT
Definition cloud.h:353
StoragePolicy storage_policy
are the functions (WorldObjects) stored in the cloud or only pointers to them
Definition cloud.h:381
std::is_constructible< T, World & > is_world_constructible
Definition cloud.h:966
valueT try_get_local_batch_bytes(const keyT record) const
bytes of a batch record held by this rank, empty if it holds none
Definition cloud.h:1024
std::atomic< long > target_replication_time
Definition cloud.h:950
nlohmann::json gather_timings(World &universe) const
Definition cloud.h:573
friend std::string to_string(const StoragePolicy sp)
Definition cloud.h:373
recordlistT store_other(madness::World &world, const T &source)
Definition cloud.h:1116
bool is_replicated
if contents of the container are replicated
Definition cloud.h:341
void set_force_load_from_cache(bool value)
Definition cloud.h:441
decltype(std::declval< T >().cloud_store(std::declval< World & >(), std::declval< Cloud & >())) member_cloud_serialize_t
Definition cloud.h:400
void replicate(const std::size_t chunk_size=INT_MAX)
Definition cloud.h:879
recordlistT local_list_of_container_keys
Definition cloud.h:394
DistributionType cloud_replication_policy
cloud is a container: replication policy for the cloud container: distributed, node-replicated,...
Definition cloud.h:384
void print_size(World &universe)
Definition cloud.h:478
std::unique_ptr< BatchTransport > batch_transport_
constructed after batch_container so it is destroyed first, as WorldObject lifetimes require
Definition cloud.h:392
std::atomic< long > reading_time
Definition cloud.h:944
ProcessID batch_owner(const keyT record) const
the owner of a batch record; a pmap lookup, no communication
Definition cloud.h:1006
void clear_timings()
Definition cloud.h:685
recordlistT store_tuple(World &world, const std::tuple< Ts... > &input)
store a tuple in multiple records
Definition cloud.h:1247
void set_debug(bool value)
Definition cloud.h:426
void register_batch_owner(const keyT record, const ProcessID owner)
Register the owner of a batch record; local map insert, no communication.
Definition cloud.h:783
nlohmann::json gather_memory_statistics(World &universe) const
get size of the cloud container
Definition cloud.h:520
std::enable_if<!is_vector< T >::value, T >::type do_load(World &world, recordlistT &recordlist) const
Definition cloud.h:1181
T load(madness::World &world, const recordlistT recordlist) const
load a single object from the cloud, recordlist is kept unchanged
Definition cloud.h:739
madness::WorldContainer< keyT, valueT > batch_container
Definition cloud.h:390
std::shared_ptr< CloudOwnerPmap< keyT > > batch_pmap
Definition cloud.h:389
void cache(madness::World &world, const T &obj, const keyT &record) const
Definition cloud.h:984
void set_fence(bool value)
Definition cloud.h:437
std::atomic< long > cache_reads
Definition cloud.h:955
std::pair< const unsigned char *, std::size_t > try_get_local_batch_ptr(const keyT record) const
stable pointer and size of a local batch record, {nullptr,0} if this rank holds none
Definition cloud.h:1036
std::vector< Function< T, NDIM > > fetch_batch_p2p(madness::World &subworld, const keyT record, const bool cache_result=false) const
fetch a batch stored by store_batch; resolves without MPI when this rank owns it
Definition cloud.h:1086
friend std::ostream & operator<<(std::ostream &os, const StoragePolicy &sp)
Definition cloud.h:364
static void print_timings(const nlohmann::json timings)
Definition cloud.h:616
bool is_rank_replicated
if every rank holds its own copy of the container (see replicate())
Definition cloud.h:342
std::list< WorldObjectBase * > world_object_base_list
Definition cloud.h:397
static void print_memory_statistics(const nlohmann::json stats)
Definition cloud.h:642
DistributionType get_replication_policy() const
is the cloud container replicated: per rank, per node, or distributed
Definition cloud.h:453
Cloud(madness::World &universe)
Definition cloud.h:408
cacheT cached_objects
Definition cloud.h:393
void clear_cache(World &subworld)
Definition cloud.h:670
keyT store_batch(madness::World &world, const std::vector< Function< T, NDIM > > &batch, const ProcessID owner, const keyT record, const bool fence=true)
Store a batch of functions as one owner-pinned record.
Definition cloud.h:797
std::atomic< long > copy_time
Definition cloud.h:949
StoragePolicy get_storing_policy() const
storing policy refers to storing functions or pointers to functions
Definition cloud.h:474
std::vector< Function< T, NDIM > > deserialize_batch_p2p(madness::World &subworld, Future< batch_bytesT > fut, const keyT record, const bool cache_result=false) const
turn the bytes of a p2p transfer into the batch of functions
Definition cloud.h:1058
std::atomic< long > batch_deserialize_time
deserializing the bytes, microseconds
Definition cloud.h:947
void replicate_according_to_policy(const std::size_t chunk_size=INT_MAX)
Definition cloud.h:842
T load_from_cache(madness::World &world, const keyT &record) const
load an object from the cache, record is unchanged
Definition cloud.h:990
madness::meta::is_detected< member_cloud_serialize_t, T > has_cloud_serialize
Definition cloud.h:403
std::vector< unsigned char > valueT
Definition cloud.h:351
void set_replication_policy(const DistributionType value)
is the cloud container replicated: per rank, per node, or distributed
Definition cloud.h:446
void distribute_targets(const DistributionType dt=Distributed)
distribute/node/rank replicate the targets of all world objects stored in the cloud
Definition cloud.h:722
bool dofence
fences after load/store
Definition cloud.h:343
void print_timings(World &universe) const
backwards compatibility
Definition cloud.h:612
bool validate_replication_policy() const
Definition cloud.h:457
T allocator(World &world) const
Definition cloud.h:1107
std::atomic< long > writing_time
Definition cloud.h:952
void set_storing_policy(const StoragePolicy value)
storing policy refers to storing functions or pointers to functions
Definition cloud.h:469
madness::WorldContainer< keyT, valueT > container
Definition cloud.h:386
T forward_load(madness::World &world, recordlistT &recordlist) const
load a single object from the cloud, recordlist is consumed while loading elements
Definition cloud.h:761
StoragePolicy
Definition cloud.h:355
@ StoreFunctionPointer
Definition cloud.h:358
@ StoreFunction
Definition cloud.h:356
std::map< keyT, cached_objT > cacheT
Definition cloud.h:352
T consuming_load(madness::World &world, recordlistT &recordlist) const
similar to load, but will consume the recordlist
Definition cloud.h:752
std::atomic< long > batch_store_time
store_batch wall time, microseconds
Definition cloud.h:945
std::atomic< long > cache_stores
Definition cloud.h:956
FunctionImpl holds all Function state to facilitate shallow copy semantics.
Definition funcimpl.h:982
A multiresolution adaptive numerical function.
Definition mra.h:144
A future is a possibly yet unevaluated value.
Definition future.h:370
T & get(bool dowork=true) &
Gets the value, waiting if necessary.
Definition future.h:571
A tensor is a multidimensional array.
Definition tensor.h:318
Makes a distributed container with specified attributes.
Definition worlddc.h:1299
bool find(accessor &acc, const keyT &key)
Write access to LOCAL value by key. Returns true if found, false otherwise (always false for remote).
Definition worlddc.h:1466
std::size_t size() const
Returns the number of local entries (no communication)
Definition worlddc.h:1614
implT::const_accessor const_accessor
Definition worlddc.h:1309
Interface to be provided by any process map.
Definition worlddc.h:125
virtual void print() const
Definition worlddc.h:142
void max(T *buf, size_t nelem)
Inplace global max while still processing AM & tasks.
Definition worldgop.h:902
void fence(bool debug=false)
Synchronizes all processes in communicator AND globally ensures no pending AM or tasks.
Definition worldgop.cc:177
void min(T *buf, size_t nelem)
Inplace global min while still processing AM & tasks.
Definition worldgop.h:896
void sum(T *buf, size_t nelem)
Inplace global sum while still processing AM & tasks.
Definition worldgop.h:890
void Bcast(T *buffer, int count, int root) const
MPI broadcast an array of count elements.
Definition worldmpi.h:416
Implements most parts of a globally addressable object (via unique ID).
Definition world_object.h:491
void process_pending()
To be called from derived constructor to process pending messages.
Definition world_object.h:787
detail::task_result_type< memfnT >::futureT send(ProcessID dest, memfnT memfn) const
Definition world_object.h:858
A parallel world class.
Definition world.h:134
ProcessID rank() const
Returns the process rank in this World (same as MPI_Comm_rank()).
Definition world.h:344
static void await(SafeMPI::Request &request, bool dowork=true)
Wait for a MPI request to complete.
Definition world.h:558
WorldMpiInterface & mpi
MPI interface.
Definition world.h:213
ProcessID size() const
Returns the number of processes in this World (same as MPI_Comm_size()).
Definition world.h:354
unsigned long id() const
Definition world.h:324
WorldGopInterface & gop
Global operations.
Definition world.h:216
void set_dofence(bool dofence)
Set the flag for fencing around a read/write operation.
Definition parallel_archive.h:302
Definition parallel_dc_archive.h:60
Definition parallel_dc_archive.h:14
long keyT
Definition parallel_dc_archive.h:16
An archive for storing local or parallel data, wrapping a BinaryFstreamInputArchive.
Definition parallel_archive.h:366
An archive for storing local or parallel data wrapping a BinaryFstreamOutputArchive.
Definition parallel_archive.h:321
Wraps an archive around an STL vector for input.
Definition vector_archive.h:101
char * p(char *buf, const char *name, int k, int initial_level, double thresh, int order)
Definition derivatives.cc:72
Tensor< typename Tensor< T >::scalar_type > arg(const Tensor< T > &t)
Return a new tensor holding the argument of each element of t (complex types only)
Definition tensor.h:2757
static const double v
Definition hatom_sf_dirac.cc:20
static double u(double r, double c)
Definition he.cc:20
#define MADNESS_CHECK(condition)
Check a condition — even in a release build the condition is always evaluated so it can have side eff...
Definition madness_exception.h:182
#define MADNESS_EXCEPTION(msg, value)
Macro for throwing a MADNESS exception.
Definition madness_exception.h:119
#define MADNESS_CHECK_THROW(condition, msg)
Check a condition — even in a release build the condition is always evaluated so it can have side eff...
Definition madness_exception.h:207
typename detail::detector< nonesuch, void, Op, Args... >::value_t is_detected
Definition meta.h:169
Definition array_addons.h:50
Namespace for all elements and tools of MADNESS.
Definition DFConvergence.h:9
void hash_range(hashT &seed, It first, It last)
Definition worldhash.h:281
double get_rss_usage_in_GB()
Definition ranks_and_hosts.cpp:18
static double cpu_time()
Returns the cpu time in seconds relative to an arbitrary origin.
Definition timers.h:128
DistributionType
some introspection of how data is distributed
Definition worlddc.h:84
@ NodeReplicated
even if there are several ranks per node
Definition worlddc.h:87
@ Distributed
no replication of the container, the container is distributed over the world
Definition worlddc.h:85
@ RankReplicated
replicate the container over all world ranks
Definition worlddc.h:86
std::vector< unsigned char > batch_bytesT
Definition cloud.h:204
void hash_combine(hashT &seed, const T &v)
Combine hash values.
Definition worldhash.h:261
static class madness::twoscale_cache_class cache[kmax+1]
DistributionType validate_distribution_type(const dcT &dc)
check distribution type of WorldContainer – global communication
Definition worlddc.h:347
void print(const T &t, const Ts &... ts)
Print items to std::cout (items separated by spaces) and terminate with a new line.
Definition print.h:227
NDIM const Function< R, NDIM > & g
Definition mra.h:2668
double wall_time()
Returns the wall time in seconds relative to an arbitrary origin.
Definition timers.cc:48
std::string type(const PairType &n)
Definition PNOParameters.h:18
vector< functionT > vecfuncT
Definition corepotential.cc:58
static bool print_timings
Definition SCF.cc:108
std::string name(const FuncType &type, const int ex=-1)
Definition ccpairfunction.h:28
madness::hashT hash_value(const std::array< T, N > &a)
Hash std::array with madness hash.
Definition array_addons.h:78
madness::archive::ContainerRecordOutputArchive::keyT batch_keyT
Definition cloud.h:201
Definition mraimpl.h:53
Definition test_dc.cc:47
Definition hatom_sf_dirac.cc:91
Definition test_ccpairfunction.cc:22
bool not_found
owner reported it does not hold the record
Definition cloud.h:279
ProcessID owner
Definition cloud.h:277
int tag
Definition cloud.h:278
std::vector< SafeMPI::Request > reqs
one per chunk, posted in on_reply
Definition cloud.h:281
Future< batch_bytesT > fut
set by finish_recv
Definition cloud.h:282
batch_bytesT buf
sized in on_reply
Definition cloud.h:280
functor to distribute/rank/node-replicate a function, passed in as a pointer to WorldObjectBase
Definition cloud.h:698
DistributionType dt
Definition cloud.h:699
int operator()(WorldObjectBase *wo) const
Definition cloud.h:701
DistributeFunctor(const DistributionType dt)
Definition cloud.h:700
Definition cloud.h:969
std::atomic< long > & rtime
Definition cloud.h:972
World & world
Definition cloud.h:970
double wall0
Definition cloud.h:971
~cloudtimer()
Definition cloud.h:976
cloudtimer(World &world, std::atomic< long > &readtime)
Definition cloud.h:974
Definition cloud.h:958
Definition cloud.h:961
Definition cloud.h:70
static keyT compute_record(const std::vector< T > &arg)
Definition cloud.h:130
keyT pop_front_and_return()
Definition cloud.h:89
static keyT compute_record(const Function< T, NDIM > &arg)
Definition cloud.h:115
Recordlist(const Recordlist &other)
Definition cloud.h:77
Recordlist(const keyT &key)
Definition cloud.h:75
static keyT compute_record(const std::shared_ptr< T > &arg)
Definition cloud.h:136
static keyT compute_record(const T &arg)
Definition cloud.h:139
std::size_t size() const
Definition cloud.h:95
madness::meta::is_detected< member_id_t, T > has_member_id
Definition cloud.h:105
decltype(std::declval< T >().id()) member_id_t
Definition cloud.h:102
friend std::ostream & operator<<(std::ostream &os, const Recordlist &arg)
Definition cloud.h:153
Recordlist & operator+=(const Recordlist &list2)
Definition cloud.h:79
static keyT compute_record(const WorldContainer< keyQ, valueT > &arg)
Definition cloud.h:121
decltype(std::declval< T >().hash()) member_hash_t
Definition cloud.h:109
static keyT compute_record(const std::shared_ptr< WorldContainer< keyQ, valueT > > &arg)
Definition cloud.h:124
std::list< keyT > list
Definition cloud.h:71
Recordlist & operator+=(const keyT &key)
Definition cloud.h:84
static keyT compute_record(const FunctionImpl< T, NDIM > *arg)
Definition cloud.h:118
static keyT compute_record(const std::shared_ptr< madness::FunctionImpl< T, NDIM > > &arg)
Definition cloud.h:127
static keyT compute_record(const Tensor< T > &arg)
Definition cloud.h:133
madness::meta::is_detected< member_hash_t, T > has_member_hash
Definition cloud.h:112
Recordlist()
Definition cloud.h:73
Base class for WorldObject.
Definition world_object.h:345
World & get_world() const
Definition world_object.h:446
class to temporarily redirect output to cout
Definition print.h:300
Definition mra.h:3042
static const char * value()
Definition cloud.h:44
static const char * value()
Definition cloud.h:46
static const char * value()
Definition cloud.h:48
static const char * value()
Definition cloud.h:50
static const char * value()
Definition cloud.h:52
static const char * value()
Definition cloud.h:54
static const char * value()
Definition cloud.h:57
static const char * value()
Definition cloud.h:59
static const char * value()
Definition cloud.h:61
static const char * value()
Definition cloud.h:63
static const char * value()
Definition cloud.h:65
static const char * value()
Definition cloud.h:67
A utility to get the name of a type as a string from chatGPT.
Definition cloud.h:39
static const char * value()
Definition cloud.h:40
#define MPI_BYTE
Definition stubmpi.h:77
double source(const coordT &r)
Definition testperiodic.cc:48
static madness::WorldMemInfo stats
Definition worldmem.cc:64
int ProcessID
Used to clearly identify process number/rank.
Definition worldtypes.h:43
FLOAT target(const FLOAT &x)
Definition y.cc:295