40#ifndef SRC_MADNESS_MRA_MACROTASKQ_H_
41#define SRC_MADNESS_MRA_MACROTASKQ_H_
79 gaxpy(1.0, x, 1.0,
true);
84 void gaxpy(
const double a,
const T& right,
double b,
const bool fence=
true) {
91 template<
typename Archive>
113template<
typename T=
double>
131 void set_impl(
const std::shared_ptr<implT>& newimpl) {
140 void gaxpy(
const double a,
const T& right,
double b,
const bool fence=
true) {
141 impl->gaxpy(
a,right,
b,fence);
144 template<
typename Archive>
156 return impl->get_local();
164 std::vector<ScalarResult<T>>
v;
217template<
typename>
struct is_tuple : std::false_type { };
218template<
typename ...T>
struct is_tuple<
std::tuple<T...>> : std::true_type { };
221template<
typename tupleT, std::
size_t I>
224 typedef decay_tuple <tupleT> argtupleT;
226 if constexpr(
I >= std::tuple_size_v<tupleT>) {
230 using typeT =
typename std::tuple_element<I, argtupleT>::type;
231 if constexpr (not is_valid_task_result_v<typeT>) {
235 return check_tuple_is_valid_task_result<tupleT,I+1>();
245 left.
gaxpy(
a, right,
b, fence);
248template <
class Archive,
typename T>
251 bool exists=(ptr) ?
true :
false;
253 if (exists) ar & ptr->id();
258template <
class Archive,
typename T>
270 MADNESS_EXCEPTION(
"ScalarResultImpl: remote operation attempting to use a locally uninitialized object",0);
273 MADNESS_EXCEPTION(
"ScalarResultImpl<T> operation attempting to use an unregistered object",0);
315 if (
name==
"default") {
319 }
else if (
name==
"node_replicated_target") {
323 }
else if (
name==
"small_memory") {
327 }
else if (
name==
"small_memory_owner") {
334 }
else if (
name==
"large_memory") {
339 std::string msg=
"MacroTaskQFactory::preset: unknown preset "+
name;
346 return {
"default",
"node_replicated_target",
"small_memory",
"small_memory_owner",
"large_memory"};
351 std::vector<MacroTaskInfo> result;
368 }
else if (store_pointer_in_cloud) {
374 if (not good) std::cout << *this ;
382 std::string msg=
"expected 3 policies, got "+std::to_string(
vec.size());
385 auto remove_quotes = [](
const std::string& s) {
386 std::string result=s;
387 if (s.size()>=2 and s.front()==
'"' and s.back()==
'"') {
388 result=s.substr(1,s.size()-2);
393 std::string sstorage=remove_quotes(
vec[0]);
398 std::string msg=
"unknown storage policy: "+sstorage;
403 std::string scloud=remove_quotes(
vec[1]);
408 std::string msg=
"unknown cloud distribution policy: "+scloud;
413 std::string sptrtarget=remove_quotes(
vec[2]);
418 std::string msg=
"unknown ptr target distribution policy: "+sptrtarget;
439 std::ostringstream os;
449 std::string msg=
"unknown policy: "+policy;
465template<
typename T=
double>
477 typedef std::vector<std::shared_ptr<MacroTaskBase> >
taskqT;
508 printf(
"this is task with priority %4.1f\n",
priority);
511 print(
"nothing to print");
514 std::stringstream ss;
536template<
typename macrotaskT>
614 std::shared_ptr< WorldDCPmapInterface< Key<1> > >
pmap1;
615 std::shared_ptr< WorldDCPmapInterface< Key<2> > >
pmap2;
616 std::shared_ptr< WorldDCPmapInterface< Key<3> > >
pmap3;
617 std::shared_ptr< WorldDCPmapInterface< Key<4> > >
pmap4;
618 std::shared_ptr< WorldDCPmapInterface< Key<5> > >
pmap5;
619 std::shared_ptr< WorldDCPmapInterface< Key<6> > >
pmap6;
673 std::shared_ptr<World> all_worlds;
674 all_worlds.reset(
new World(comm));
687 std::shared_ptr<World> node_world(
new World(comm));
717 if (need_replication_of_target) {
720 loop_types<Cloud::DistributeFunctor, double, float, double_complex, float_complex>(std::tuple<DistributionType>(dt),wo);
747 const bool all_tasks_owned = (not
taskq.empty()) and std::all_of(
taskq.begin(),
taskq.end(),
748 [](
const std::shared_ptr<MacroTaskBase>& t) {
return t->get_owner_slot() >= 0; });
751 for (
long element=0; element<long(
taskq.size()); ++element) {
752 std::shared_ptr<MacroTaskBase>
task=
taskq[element];
753 if (not
task->is_owned_by(requester_slot))
continue;
758 tasktime+=(cpu1-cpu0);
759 if (
printdebug()) printf(
"completed task %3ld after %6.1fs at time %6.1fs\n",element,cpu1-cpu0,
wall_time());
762 printf(
"rank %3d (subworld %3lu) finished its task queue at time %6.1fs (own tasktime %4.1fs)\n",
771 if (element<0)
break;
772 std::shared_ptr<MacroTaskBase>
task=
taskq[element];
779 tasktime+=(cpu1-cpu0);
780 if (
printdebug()) printf(
"completed task %3ld after %6.1fs at time %6.1fs\n",element,cpu1-cpu0,
wall_time());
783 const std::size_t ntask=
taskq.size();
785 auto in_percentile = [&ntask](
const long element) {
786 return std::floor(element/(0.1*(ntask+1)));
788 auto is_first_in_percentile = [&](
const long element) {
789 return (in_percentile(element)!=in_percentile(element-1));
792 std::cout << int(in_percentile(element)*10) <<
" " << std::flush;
821 const bool want_node_reduction = (not
taskq.empty())
822 and
taskq.front()->wants_node_local_reduction();
827 if (nodeworld and nodeworld->
size() ==
universe.
size()) nodeworld =
nullptr;
829 for (
auto& t :
taskq) t->finalize_stage1(subworld, nodeworld,
cloud);
830 if (nodeworld) nodeworld->
gop.
fence();
831 for (
auto& t :
taskq) t->finalize_stage2(subworld, nodeworld,
cloud);
839 print(
"number of tasks in taskq",
taskq.size());
840 print(
"redirecting output to files task.#####");
876 print(
"all tasks complete");
881 printf(
"completed taskqueue after %4.1fs at time %4.1fs\n", cpu11 - cpu00,
wall_time());
882 printf(
" total cpu time / per world %4.1fs %4.1fs\n", tasktime, tasktime /
universe.
size());
912 for (
const auto& t : vtask) {
922 print(
"total number of tasks: ",
taskq.size());
923 print(
" task batch priority status");
924 for (
const auto& t :
taskq) t->print_me_as_table();
937 if (subworld.
rank()==0) {
951 auto is_Waiting = [](
const std::shared_ptr<MacroTaskBase>& mtb_ptr) {
return mtb_ptr->is_waiting();};
952 auto it=std::find_if(
taskq.begin(),
taskq.end(),is_Waiting);
953 if (it!=
taskq.end()) {
954 it->get()->set_running();
955 long element=it-
taskq.begin();
970 taskq[task_number]->set_complete();
990template<
typename taskT>
1016 template<
typename Q>
1018 decltype(std::declval<Q&>().prepare_owner_assignment(
1019 std::declval<const MacroTaskPartitioner::partitionT&>(), 0
L));
1020 template<
typename Q>
1022 madness::meta::is_detected_v<has_prepare_owner_assignment_t, Q>;
1025 template<
typename Q,
typename ArgTuple>
1027 decltype(std::declval<Q&>().store_batches(
1028 std::declval<World&>(), std::declval<World&>(), std::declval<Cloud&>(),
1029 std::declval<const ArgTuple&>(), 0
L));
1030 template<
typename Q,
typename ArgTuple>
1032 madness::meta::is_detected_v<has_store_batches_t, Q, ArgTuple>;
1035 template<
typename Q,
typename VecT>
1037 decltype(std::declval<Q&>().sym_pipeline_advance(
1038 std::declval<World&>(), std::declval<const VecT&>(),
1039 std::declval<const Batch_1D&>(), std::declval<const Batch_1D&>(),
true));
1040 template<
typename Q,
typename VecT>
1042 madness::meta::is_detected_v<has_sym_pipeline_advance_t, Q, VecT>;
1045 template<
typename Q,
typename ResT>
1047 decltype(std::declval<Q&>().accumulate_locally(
1048 std::declval<World&>(), std::declval<const ResT&>()));
1049 template<
typename Q,
typename ResT>
1051 madness::meta::is_detected_v<has_accumulate_locally_t, Q, ResT>;
1054 template<
typename Q>
1056 decltype(std::declval<const Q&>().wants_node_local_reduction());
1057 template<
typename Q>
1059 madness::meta::is_detected_v<has_wants_node_local_reduction_t, Q>;
1062 template<
typename Q>
1064 decltype(std::declval<Q&>().finalize_stage1(std::declval<World&>(),
1065 std::declval<World*>()));
1066 template<
typename Q>
1068 madness::meta::is_detected_v<has_finalize_stage1_t, Q>;
1071 template<
typename Q,
typename ResT>
1073 decltype(std::declval<Q&>().finalize_stage2(std::declval<World&>(),
1074 std::declval<World*>(),
1075 std::declval<ResT&>()));
1076 template<
typename Q,
typename ResT>
1078 madness::meta::is_detected_v<has_finalize_stage2_t, Q, ResT>;
1110 if (this->taskq_ptr==0) {
1115 if (
debug) this->taskq_ptr->set_printlevel(20);
1118 this->taskq_ptr->cloud.set_storing_policy(cloud_storage_policy);
1119 this->taskq_ptr->cloud.set_replication_policy(this->taskq_ptr->get_policy().cloud_distribution_policy);
1137 template<
typename ... Ts>
1140 auto argtuple = std::tie(args...);
1141 static_assert(std::is_same<
decltype(argtuple),
argtupleT>::value,
"type or number of arguments incorrect");
1144 auto partitioner=
task.partitioner;
1146 partitioner->set_nsubworld(
world.
size());
1147 partitionT partition = partitioner->partition_tasks(argtuple);
1150 if constexpr (has_prepare_owner_assignment_v<taskT>)
1151 task.prepare_owner_assignment(partition,
taskq_ptr->get_nsubworld());
1157 if constexpr (has_store_batches_v<taskT, argtupleT>)
1165 for (
const auto& batch_prio : partition) {
1166 const long owner_slot =
task.owner_hint(batch_prio.first,
taskq_ptr->get_nsubworld());
1168 std::shared_ptr<MacroTaskBase>(
new MacroTaskInternal(
task, batch_prio, inputrecords, outputrecords, owner_slot)));
1181 static_assert(check_tuple_is_valid_task_result<resultT,0>(),
1182 "tuple has invalid result type in prepare_output_records");
1184 static_assert(is_valid_task_result_v<resultT>,
"unknown result type in prepare_output_records");
1187 if (
debug)
print(
"storing pointers to output in cloud");
1189 auto store_output_records = [&](
const auto& result) {
1191 typedef std::decay_t<
decltype(result)> argT;
1193 outputrecords += cloud.
store(
world, result.get_impl().get());
1197 outputrecords += cloud.
store(
world, result.get_impl());
1201 std::vector<std::shared_ptr<typename argT::value_type::implT>>
v;
1202 for (
const auto& ptr : result)
v.push_back(ptr.get_impl());
1210 return outputrecords;
1216 std::apply([&](
auto &&... args) {
1217 (( outputrecords+=store_output_records(args) ), ...);
1220 outputrecords=store_output_records(result);
1222 return outputrecords;
1235 if (
task.name==
"unknown_task")
return typeid(
task).
name();
1244 static_assert(check_tuple_is_valid_task_result<resultT,0>(),
1245 "tuple has invalid result type in prepare_output_records");
1247 static_assert(is_valid_task_result_v<resultT>,
"unknown result type in prepare_output_records");
1249 this->task.batch=batch_prio.first;
1256 print(
"this is task",
get_name(),
"with batch",
task.batch,
"priority",this->get_priority());
1260 std::stringstream ss;
1262 std::size_t namesize=std::min(std::size_t(28),
name.size());
1263 name += std::string(28-namesize,
' ');
1265 std::stringstream ssbatch;
1266 ssbatch <<
task.batch;
1267 std::string strbatch=ssbatch.str();
1268 int nspaces=std::max(
int(0),35-
int(ssbatch.str().size()));
1269 strbatch+=std::string(nspaces,
' ');
1272 << std::setw(10) << strbatch
1278 template<
typename resultT1, std::
size_t I=0>
1279 typename std::enable_if<is_tuple<resultT1>::value,
void>
::type
1281 if constexpr(I < std::tuple_size_v<resultT1>) {
1282 using elementT =
typename std::tuple_element<I, resultT>::type;
1283 auto element_final=std::get<I>(final_result);
1284 auto element_tmp=std::get<I>(tmp_result);
1285 accumulate_into_final_result<elementT>(subworld, element_final, element_tmp, argtuple);
1286 accumulate_into_final_result<resultT1,I+1>(subworld, final_result, tmp_result, argtuple);
1291 template<
typename resultT1>
1292 typename std::enable_if<not is_tuple<resultT1>::value,
void>
::type
1297 result_tmp.change_tree_state(operating_state);
1298 gaxpy(1.0,result,1.0, result_tmp);
1305 gaxpy(1.0,result,1.0,result_tmp,
false);
1309 gaxpy(1.0, result, 1.0, result_tmp.get_local(),
false);
1313 std::size_t sz=result.size();
1314 for (
size_t i=0; i<sz; ++i) {
1315 gaxpy(1.0, result[i], 1.0, result_tmp[i].get_local(),
false);
1326 if constexpr (has_wants_node_local_reduction_v<taskT>)
return task.wants_node_local_reduction();
1331 if (not
task.accumulates_own_output())
return;
1332 if constexpr (has_finalize_stage1_v<taskT>)
task.finalize_stage1(subworld, nodeworld);
1336 if (not
task.accumulates_own_output())
return;
1338 if constexpr (has_finalize_stage2_v<taskT, resultT>)
1339 task.finalize_stage2(subworld, nodeworld, result_universe);
1350 if (owner_slot < 0)
return {-1,
nullptr};
1351 for (
long next = element + 1; next < long(taskq.size()); ++next) {
1353 auto next_task = std::dynamic_pointer_cast<MacroTaskInternal>(taskq[next]);
1354 if (next_task)
return {next, next_task.get()};
1356 return {-1,
nullptr};
1372 const bool has_next = (next_ptr !=
nullptr) and (next_ptr->task.batch.input.size() > 1);
1373 const Batch_1D next_col = has_next ? next_ptr->task.batch.input[0] :
Batch_1D();
1374 const Batch_1D next_row = has_next ? next_ptr->task.batch.input[1] :
Batch_1D();
1376 if constexpr (std::tuple_size<argtupleT>::value >= 3) {
1377 using ketT = std::decay_t<std::tuple_element_t<2, argtupleT>>;
1378 if constexpr (has_sym_pipeline_advance_v<taskT, ketT>)
1379 task.sym_pipeline_advance(subworld, std::get<2>(argtuple),
1380 next_col, next_row, has_next);
1395 const bool need_auto_copy =
1398 and not
task.handles_own_data_movement();
1399 if (need_auto_copy) {
1408 auto copi = [&](
auto&
arg) {
1409 typedef std::decay_t<
decltype(
arg)> argT;
1421 print(
"copied coefficients for task",
get_name(),
"in",cpu1-cpu0,
"seconds");
1437 argtupleT batched_argtuple =
task.batch.copy_input_batch(argtuple);
1439 task.subworld_ptr=&subworld;
1441 task.cloud_ptr=&cloud;
1446 print(
"starting task no",element,
", '",
get_name(),
"', in subworld",subworld.id(),
"at time",
wall_time());
1448 resultT result_batch = std::apply(
task, batched_argtuple);
1450 constexpr std::size_t
bufsize=256;
1452 std::snprintf(buffer,
bufsize,
"completed task %3ld after %6.1fs at time %6.1fs\n",element,cpu1-cpu0,
wall_time());
1453 print(std::string(buffer));
1456 auto insert_batch = [&](
auto& element1,
auto& element2) {
1457 typedef std::decay_t<
decltype(element1)> decay_type;;
1459 element1=
task.batch.insert_result_batch(element1,element2);
1461 std::swap(element1,element2);
1464 resultT result_subworld=
task.allocator(subworld,argtuple);
1468 insert_batch(result_subworld,result_batch);
1474 if (
task.accumulates_own_output()) {
1475 if constexpr (has_accumulate_locally_v<taskT, resultT>)
1476 task.accumulate_locally(subworld, result_subworld);
1479 accumulate_into_final_result<resultT>(subworld, result_universe, result_subworld, argtuple);
1482 }
catch (std::exception&
e) {
1486 print(
"failing task no",element,
"in subworld",subworld.
id(),
"at time",
wall_time());
1488 print(
"RSS at failure (current resident, GB):", rss_at_fail);
1494 print(
"failing task no",element,
"in subworld",subworld.
id(),
"at time",
wall_time());
1496 print(
"RSS at failure (current resident, GB):", rss_at_fail);
1508 template<
typename T, std::
size_t NDIM>
1515 template<
typename T, std::
size_t NDIM>
1517 std::vector<Function<T,NDIM>> vresult;
1518 vresult.resize(v_impl.size());
1523 template<
typename T>
1528 template<
typename T>
1530 std::vector<ScalarResult<T>> vresult(v_sr_impl.size());
1531 for (
size_t i=0; i<v_sr_impl.size(); ++i) {
1532 vresult[i].set_impl(v_sr_impl[i]);
1548 auto doit = [&](
auto& element) {
1549 typedef std::decay_t<
decltype(element)> elementT;
1553 typedef typename elementT::value_type::implT implT;
1554 auto ptr_element = cloud.
consuming_load<std::vector<std::shared_ptr<implT>>>(
1555 subworld, outputrecords1);
1559 typedef typename elementT::implT implT;
1560 auto ptr_element = cloud.
consuming_load<std::shared_ptr<implT>>(subworld, outputrecords1);
1564 typedef typename elementT::value_type ScalarResultT;
1565 typedef typename ScalarResultT::implT implT;
1566 typedef std::vector<std::shared_ptr<implT>> vptrT;
1567 auto ptr_element = cloud.
consuming_load<vptrT>(subworld, outputrecords1);
1573 auto ptr_element = cloud.
consuming_load<std::shared_ptr<typename elementT::implT>>(subworld, outputrecords1);
1581 static_assert(check_tuple_is_valid_task_result<resultT, 0>(),
1582 "invalid tuple task result -- must be vectors of functions");
Wrapper around MPI_Comm. Has a shallow copy constructor; use Create(Get_group()) for deep copy.
Definition safempi.h:497
static const int SHARED_SPLIT_TYPE
Definition safempi.h:657
Intracomm Split_type(int Type, int Key=0) const
Definition safempi.h:674
Intracomm Split(int Color, int Key=0) const
Definition safempi.h:642
Definition macrotaskpartitioner.h:55
a batch consists of a 2D-input batch and a 1D-output batch: K-batch <- (I-batch, J-batch)
Definition macrotaskpartitioner.h:124
cloud class
Definition cloud.h:338
void clear()
Definition cloud.h:675
void replicate_per_node(const std::size_t chunk_size=INT_MAX)
Definition cloud.h:858
nlohmann::json get_statistics(World &world) const
return a json object with the cloud settings and statistics
Definition cloud.h:503
recordlistT store(madness::World &world, const T &source)
Definition cloud.h:818
Recordlist< keyT > recordlistT
Definition cloud.h:352
std::atomic< long > target_replication_time
Definition cloud.h:948
nlohmann::json gather_timings(World &universe) const
Definition cloud.h:572
void replicate(const std::size_t chunk_size=INT_MAX)
Definition cloud.h:878
void print_size(World &universe)
Definition cloud.h:477
T load(madness::World &world, const recordlistT recordlist) const
load a single object from the cloud, recordlist is kept unchanged
Definition cloud.h:738
std::list< WorldObjectBase * > world_object_base_list
Definition cloud.h:396
static void print_memory_statistics(const nlohmann::json stats)
Definition cloud.h:641
DistributionType get_replication_policy() const
is the cloud container replicated: per rank, per node, or distributed
Definition cloud.h:452
void clear_cache(World &subworld)
Definition cloud.h:669
std::atomic< long > copy_time
Definition cloud.h:947
void set_replication_policy(const DistributionType value)
is the cloud container replicated: per rank, per node, or distributed
Definition cloud.h:445
void distribute_targets(const DistributionType dt=Distributed)
distribute/node/rank replicate the targets of all world objects stored in the cloud
Definition cloud.h:721
void print_timings(World &universe) const
backwards compatibility
Definition cloud.h:611
void set_storing_policy(const StoragePolicy value)
storing policy refers to storing functions or pointers to functions
Definition cloud.h:468
StoragePolicy
Definition cloud.h:354
@ StoreFunctionPointer
Definition cloud.h:357
@ StoreFunction
Definition cloud.h:355
T consuming_load(madness::World &world, recordlistT &recordlist) const
similar to load, but will consume the recordlist
Definition cloud.h:751
static void set_default_pmap(World &world)
Definition mraimpl.h:3715
static std::shared_ptr< WorldDCPmapInterface< Key< NDIM > > > & get_pmap()
Returns the default process map that was last initialized via set_default_pmap()
Definition funcdefaults.h:399
static void set_pmap(const std::shared_ptr< WorldDCPmapInterface< Key< NDIM > > > &value)
Sets the default process map (does not redistribute existing functions)
Definition funcdefaults.h:430
FunctionImpl holds all Function state to facilitate shallow copy semantics.
Definition funcimpl.h:970
A multiresolution adaptive numerical function.
Definition mra.h:144
void set_impl(const std::shared_ptr< FunctionImpl< T, NDIM > > &impl)
Replace current FunctionImpl with provided new one.
Definition mra.h:731
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
base class
Definition macrotaskq.h:474
virtual void finalize_stage1(World &, World *, Cloud &)
Definition macrotaskq.h:504
void set_running()
Definition macrotaskq.h:487
bool is_owned_by(const long slot) const
an unowned task may be run by any subworld, an owned one only by its own
Definition macrotaskq.h:524
virtual ~MacroTaskBase()
Definition macrotaskq.h:480
void set_waiting()
Definition macrotaskq.h:488
MacroTaskBase()
Definition macrotaskq.h:479
virtual void print_me(std::string s="") const
Definition macrotaskq.h:507
std::string print_priority_and_status_to_string() const
Definition macrotaskq.h:513
virtual bool wants_node_local_reduction() const
Finalize, for tasks that accumulate their own output.
Definition macrotaskq.h:503
void set_complete()
Definition macrotaskq.h:486
double priority
Definition macrotaskq.h:482
bool is_complete() const
Definition macrotaskq.h:490
Status
Definition macrotaskq.h:484
@ Complete
Definition macrotaskq.h:484
@ Running
Definition macrotaskq.h:484
@ Unknown
Definition macrotaskq.h:484
@ Waiting
Definition macrotaskq.h:484
void set_owner_slot(const long owner)
Definition macrotaskq.h:522
bool is_running() const
Definition macrotaskq.h:491
long get_owner_slot() const
Definition macrotaskq.h:521
enum madness::MacroTaskBase::Status stat
virtual void finalize_stage2(World &, World *, Cloud &)
Definition macrotaskq.h:505
long owner_slot
the subworld that should run this task; -1 means any
Definition macrotaskq.h:483
double get_priority() const
Definition macrotaskq.h:519
void set_priority(const double p)
Definition macrotaskq.h:520
virtual void print_me_as_table(std::string s="") const
Definition macrotaskq.h:510
virtual void run(World &world, Cloud &cloud, taskqT &taskq, const long element, const bool debug, const MacroTaskInfo policy)=0
std::vector< std::shared_ptr< MacroTaskBase > > taskqT
Definition macrotaskq.h:477
friend std::ostream & operator<<(std::ostream &os, const MacroTaskBase::Status s)
Definition macrotaskq.h:526
bool is_waiting() const
Definition macrotaskq.h:492
Definition macrotaskq.h:1604
Cloud * cloud_ptr
Definition macrotaskq.h:1610
MacroTaskOperationBase()
Definition macrotaskq.h:1613
Batch batch
Definition macrotaskq.h:1606
World * subworld_ptr
Definition macrotaskq.h:1607
virtual bool accumulates_own_output() const
Definition macrotaskq.h:1624
virtual long owner_hint(const Batch &, const long) const
which subworld should run this batch; -1 leaves the choice to the queue
Definition macrotaskq.h:1620
virtual ~MacroTaskOperationBase()
Definition macrotaskq.h:1614
std::shared_ptr< MacroTaskPartitioner > partitioner
Definition macrotaskq.h:1612
virtual bool handles_own_data_movement() const
true if the task moves its operand coefficients into the subworld itself
Definition macrotaskq.h:1627
virtual void cleanup()
release whatever the task kept across its batches, called once the queue has run
Definition macrotaskq.h:1636
std::string name
Definition macrotaskq.h:1611
partition one (two) vectors into 1D (2D) batches.
Definition macrotaskpartitioner.h:182
std::list< std::pair< Batch, double > > partitionT
Definition macrotaskpartitioner.h:186
Factory for the MacroTaskQ.
Definition macrotaskq.h:550
MacroTaskQFactory & set_storage_policy(const MacroTaskInfo::StoragePolicy sp)
Definition macrotaskq.h:579
World & world
Definition macrotaskq.h:553
MacroTaskQFactory & set_policy(const MacroTaskInfo p)
Definition macrotaskq.h:574
MacroTaskQFactory & set_nworld(const long n)
Definition macrotaskq.h:560
MacroTaskQFactory & preset(const std::string name)
Definition macrotaskq.h:570
MacroTaskQFactory & set_cloud_distribution_policy(const DistributionType dp)
Definition macrotaskq.h:584
MacroTaskQFactory & set_printlevel(const long p)
Definition macrotaskq.h:565
long printlevel
Definition macrotaskq.h:552
MacroTaskQFactory & set_ptr_target_distribution_policy(const DistributionType dp)
Definition macrotaskq.h:589
long nworld
Definition macrotaskq.h:554
MacroTaskInfo policy
Definition macrotaskq.h:556
MacroTaskQFactory(World &universe)
Definition macrotaskq.h:558
Definition macrotaskq.h:598
long nsubworld
Definition macrotaskq.h:607
std::mutex taskq_mutex
Definition macrotaskq.h:605
bool printdebug() const
Definition macrotaskq.h:621
static void set_pmap(World &world)
Definition macrotaskq.h:974
bool printtimings_detail() const
Definition macrotaskq.h:624
void run_all()
Definition macrotaskq.h:834
void set_complete(const long task_number) const
scheduler is located on rank==0
Definition macrotaskq.h:963
World & universe
Definition macrotaskq.h:600
std::shared_ptr< WorldDCPmapInterface< Key< 1 > > > pmap1
set the process map for the subworld
Definition macrotaskq.h:614
std::shared_ptr< WorldDCPmapInterface< Key< 3 > > > pmap3
Definition macrotaskq.h:616
MacroTaskBase::taskqT taskq
Definition macrotaskq.h:604
MacroTaskInfo get_policy() const
Definition macrotaskq.h:633
nlohmann::json cloud_statistics
save cloud statistics after run_all()
Definition macrotaskq.h:608
void add_tasks(MacroTaskBase::taskqT &vtask)
Definition macrotaskq.h:911
long get_scheduled_task_number(World &subworld)
scheduler is located on universe.rank==0
Definition macrotaskq.h:935
void set_complete_local(const long task_number) const
scheduler is located on rank==0
Definition macrotaskq.h:968
World & get_subworld()
Definition macrotaskq.h:629
MacroTaskQ(const MacroTaskQFactory factory)
create an empty taskq and initialize the subworlds
Definition macrotaskq.h:646
static std::shared_ptr< World > create_node_world(World &universe)
a World spanning the ranks that share memory with this one
Definition macrotaskq.h:684
void drain_own_output_buffers()
Move the results of tasks that accumulated their own output into the universe result.
Definition macrotaskq.h:803
std::shared_ptr< WorldDCPmapInterface< Key< 4 > > > pmap4
Definition macrotaskq.h:617
nlohmann::json get_taskq_statistics() const
Definition macrotaskq.h:641
std::size_t size() const
Definition macrotaskq.h:983
std::shared_ptr< WorldDCPmapInterface< Key< 5 > > > pmap5
Definition macrotaskq.h:618
std::shared_ptr< World > nodeworld_ptr
node-scoped World for the two-stage finalize; created on first use, reused after
Definition macrotaskq.h:603
void print_taskq() const
Definition macrotaskq.h:918
std::shared_ptr< World > subworld_ptr
Definition macrotaskq.h:601
bool printtimings() const
Definition macrotaskq.h:623
std::shared_ptr< WorldDCPmapInterface< Key< 2 > > > pmap2
Definition macrotaskq.h:615
double execute_tasks()
Run this rank's share of the queue.
Definition macrotaskq.h:737
bool printprogress() const
Definition macrotaskq.h:622
void replicate_inputs()
run all tasks
Definition macrotaskq.h:698
std::shared_ptr< WorldDCPmapInterface< Key< 6 > > > pmap6
Definition macrotaskq.h:619
nlohmann::json get_cloud_statistics() const
Definition macrotaskq.h:637
void set_printlevel(const long p)
Definition macrotaskq.h:631
long printlevel
Definition macrotaskq.h:606
madness::Cloud cloud
Definition macrotaskq.h:628
long get_nsubworld() const
Definition macrotaskq.h:630
void add_replicated_task(const std::shared_ptr< MacroTaskBase > &task)
Definition macrotaskq.h:930
~MacroTaskQ()
Definition macrotaskq.h:664
static std::shared_ptr< World > create_worlds(World &universe, const std::size_t nsubworld)
Definition macrotaskq.h:668
nlohmann::json taskq_statistics
save taskq statistics after run_all()
Definition macrotaskq.h:609
long get_scheduled_task_number_local()
Definition macrotaskq.h:947
const MacroTaskInfo policy
storage and distribution policy
Definition macrotaskq.h:611
Definition macrotaskq.h:1226
void run(World &subworld, Cloud &cloud, MacroTaskBase::taskqT &taskq, const long element, const bool debug, const MacroTaskInfo policy) override
Definition macrotaskq.h:1426
void finalize_stage2(World &subworld, World *nodeworld, Cloud &cloud) override
Definition macrotaskq.h:1335
void finalize_stage1(World &subworld, World *nodeworld, Cloud &) override
Definition macrotaskq.h:1330
taskT task
Definition macrotaskq.h:1233
void prefetch_for_next_task(World &subworld, const argtupleT &argtuple, MacroTaskBase::taskqT &taskq, const long element)
called by the MacroTaskQ when the task is scheduled
Definition macrotaskq.h:1366
void copy_operands_into_subworld(World &subworld, Cloud &cloud, argtupleT &batched_argtuple, const MacroTaskInfo &policy, const bool debug)
Bring the operand coefficients into the subworld, unless the task does that itself.
Definition macrotaskq.h:1386
std::enable_if< notis_tuple< resultT1 >::value, void >::type accumulate_into_final_result(World &subworld, resultT1 &result, const resultT1 &result_tmp, const argtupleT &argtuple)
accumulate the result of the task into the final result living in the universe
Definition macrotaskq.h:1293
resultT get_output(World &subworld, Cloud &cloud) const
return the WorldObjects or the result functions living in the universe
Definition macrotaskq.h:1541
static ScalarResult< T > pointer2WorldObject(const std::shared_ptr< ScalarResultImpl< T > > sr_impl)
Definition macrotaskq.h:1524
decay_tuple< typename taskT::argtupleT > argtupleT
Definition macrotaskq.h:1228
static std::vector< Function< T, NDIM > > pointer2WorldObject(const std::vector< std::shared_ptr< FunctionImpl< T, NDIM > > > v_impl)
Definition macrotaskq.h:1516
void cleanup() override
Definition macrotaskq.h:1504
taskT::resultT resultT
Definition macrotaskq.h:1229
recordlistT outputrecords
Definition macrotaskq.h:1231
std::string get_name() const
Definition macrotaskq.h:1234
static Function< T, NDIM > pointer2WorldObject(const std::shared_ptr< FunctionImpl< T, NDIM > > impl)
Definition macrotaskq.h:1509
static std::vector< ScalarResult< T > > pointer2WorldObject(const std::vector< std::shared_ptr< ScalarResultImpl< T > > > v_sr_impl)
Definition macrotaskq.h:1529
void print_me_as_table(std::string s="") const override
Definition macrotaskq.h:1259
bool wants_node_local_reduction() const override
Bridge the framework's virtual finalize to the task's hooks.
Definition macrotaskq.h:1325
void print_me(std::string s="") const override
Definition macrotaskq.h:1255
std::enable_if< is_tuple< resultT1 >::value, void >::type accumulate_into_final_result(World &subworld, resultT1 &final_result, const resultT1 &tmp_result, const argtupleT &argtuple)
accumulate the result of the task into the final result living in the universe
Definition macrotaskq.h:1280
MacroTaskInternal(const taskT &task, const std::pair< Batch, double > &batch_prio, const recordlistT &inputrecords, const recordlistT &outputrecords, const long owner_slot=-1)
Definition macrotaskq.h:1239
std::pair< long, MacroTaskInternal * > find_next_owned_task(const MacroTaskBase::taskqT &taskq, const long element) const
the next task in taskq that this task's subworld also owns, or {-1, nullptr}
Definition macrotaskq.h:1347
recordlistT inputrecords
Definition macrotaskq.h:1230
Definition macrotaskq.h:991
decltype(std::declval< Q & >().sym_pipeline_advance(std::declval< World & >(), std::declval< const VecT & >(), std::declval< const Batch_1D & >(), std::declval< const Batch_1D & >(), true)) has_sym_pipeline_advance_t
true if: task.sym_pipeline_advance(subworld, ket, next_col, next_row, has_next)
Definition macrotaskq.h:1039
MacroTask & set_debug(const bool value)
Definition macrotaskq.h:1123
MacroTask(World &world, taskT &task)
constructor takes the task, but no arguments to the task
Definition macrotaskq.h:1093
MacroTask(World &world, taskT &task, const MacroTaskQFactory factory)
constructor takes task and a taskq factory for customization, immediate execution
Definition macrotaskq.h:1099
taskT::resultT resultT
Definition macrotaskq.h:1080
static constexpr bool has_finalize_stage2_v
Definition macrotaskq.h:1077
std::shared_ptr< MacroTaskQ > taskq_ptr
Definition macrotaskq.h:1088
taskT::argtupleT argtupleT
Definition macrotaskq.h:1081
bool immediate_execution
Definition macrotaskq.h:1086
decltype(std::declval< Q & >().finalize_stage1(std::declval< World & >(), std::declval< World * >())) has_finalize_stage1_t
true if: task.finalize_stage1(subworld, nodeworld)
Definition macrotaskq.h:1065
static constexpr bool has_wants_node_local_reduction_v
Definition macrotaskq.h:1058
static constexpr bool has_finalize_stage1_v
Definition macrotaskq.h:1067
World & world
Definition macrotaskq.h:1087
std::shared_ptr< MacroTaskQ > get_taskq() const
Definition macrotaskq.h:1128
resultT operator()(const Ts &... args)
this mimicks the original call to the task functor, called from the universe
Definition macrotaskq.h:1138
static constexpr bool has_accumulate_locally_v
Definition macrotaskq.h:1050
static constexpr bool has_sym_pipeline_advance_v
Definition macrotaskq.h:1041
decltype(std::declval< Q & >().accumulate_locally(std::declval< World & >(), std::declval< const ResT & >())) has_accumulate_locally_t
true if: task.accumulate_locally(subworld, result_subworld)
Definition macrotaskq.h:1048
bool debug
Definition macrotaskq.h:1085
taskT task
Definition macrotaskq.h:1084
MacroTaskPartitioner::partitionT partitionT
Definition macrotaskq.h:992
decltype(std::declval< Q & >().store_batches(std::declval< World & >(), std::declval< World & >(), std::declval< Cloud & >(), std::declval< const ArgTuple & >(), 0L)) has_store_batches_t
true if: task.store_batches(world, subworld, cloud, argtuple, nsubworld)
Definition macrotaskq.h:1029
decltype(std::declval< Q & >().finalize_stage2(std::declval< World & >(), std::declval< World * >(), std::declval< ResT & >())) has_finalize_stage2_t
true if: task.finalize_stage2(subworld, nodeworld, universe_result)
Definition macrotaskq.h:1075
Cloud::recordlistT recordlistT
Definition macrotaskq.h:1082
static constexpr bool has_store_batches_v
Definition macrotaskq.h:1031
static constexpr bool has_prepare_owner_assignment_v
Definition macrotaskq.h:1021
recordlistT prepare_output_records(Cloud &cloud, resultT &result)
store pointers to the result WorldObject in the cloud and return the recordlist
Definition macrotaskq.h:1179
decltype(std::declval< Q & >().prepare_owner_assignment(std::declval< const MacroTaskPartitioner::partitionT & >(), 0L)) has_prepare_owner_assignment_t
true if: task.prepare_owner_assignment(partition, nsubworld)
Definition macrotaskq.h:1019
MacroTask(World &world, taskT &task, std::shared_ptr< MacroTaskQ > taskq_ptr)
constructor takes the task, and a taskq, execution is not immediate
Definition macrotaskq.h:1105
decltype(std::declval< const Q & >().wants_node_local_reduction()) has_wants_node_local_reduction_t
true if: task.wants_node_local_reduction()
Definition macrotaskq.h:1056
static std::map< MemKey, MemInfo > measure_and_print(World &world)
measure the memory usage of all objects of all worlds
Definition memory_measurement.h:24
helper class for returning the result of a task, which is not a madness Function, but a simple scalar
Definition macrotaskq.h:55
void gaxpy(const double a, const T &right, double b, const bool fence=true)
accumulate, optional fence
Definition macrotaskq.h:84
void serialize(Archive &ar)
Definition macrotaskq.h:92
T get()
after completion of the taskq get the final value
Definition macrotaskq.h:97
ScalarResultImpl< T > & operator=(const T &x)
simple assignment of the scalar value
Definition macrotaskq.h:73
ScalarResultImpl(const ScalarResultImpl &other)=delete
Disable the default copy constructor.
ScalarResultImpl< T > & operator=(const ScalarResultImpl< T > &other)=delete
disable assignment operator
ScalarResultImpl< T > & operator+=(const T &x)
Definition macrotaskq.h:78
~ScalarResultImpl()
Definition macrotaskq.h:69
T value
the scalar value
Definition macrotaskq.h:110
T get_local() const
Definition macrotaskq.h:104
T value_type
Definition macrotaskq.h:57
ScalarResultImpl(World &world)
Definition macrotaskq.h:58
Definition macrotaskq.h:114
void serialize(Archive &ar)
Definition macrotaskq.h:145
ScalarResult(const std::shared_ptr< implT > &impl)
Definition macrotaskq.h:121
void gaxpy(const double a, const T &right, double b, const bool fence=true)
accumulate, optional fence
Definition macrotaskq.h:140
ScalarResultImpl< T > implT
Definition macrotaskq.h:116
T get()
after completion of the taskq get the final value
Definition macrotaskq.h:150
ScalarResult(World &world)
Definition macrotaskq.h:120
std::shared_ptr< implT > impl
Definition macrotaskq.h:117
ScalarResult & operator=(const T &x)
Definition macrotaskq.h:122
std::shared_ptr< implT > get_impl() const
Definition macrotaskq.h:127
void set_impl(const std::shared_ptr< implT > &newimpl)
Definition macrotaskq.h:131
T get_local() const
after completion of the taskq get the final value
Definition macrotaskq.h:155
uniqueidT id() const
Definition macrotaskq.h:135
void broadcast_serializable(objT &obj, ProcessID root)
Broadcast a serializable object.
Definition worldgop.h:774
void fence(bool debug=false)
Synchronizes all processes in communicator AND globally ensures no pending AM or tasks.
Definition worldgop.cc:176
bool set_forbid_fence(bool value)
Set forbid_fence flag to new value and return old value.
Definition worldgop.h:677
void sum(T *buf, size_t nelem)
Inplace global sum while still processing AM & tasks.
Definition worldgop.h:890
SafeMPI::Intracomm & comm()
Returns the associated SafeMPI communicator.
Definition worldmpi.h:286
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
detail::task_result_type< memfnT >::futureT task(ProcessID dest, memfnT memfn, const TaskAttributes &attr=TaskAttributes()) const
Sends task to derived class method returnT (this->*memfn)().
Definition world_object.h:1132
A parallel world class.
Definition world.h:134
static World * world_from_id(std::uint64_t id)
Convert a World ID to a World pointer.
Definition world.h:516
ProcessID rank() const
Returns the process rank in this World (same as MPI_Comm_rank()).
Definition world.h:344
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
std::optional< T * > ptr_from_id(uniqueidT id) const
Look up a local pointer from a world-wide unique ID.
Definition world.h:440
Class for unique global IDs.
Definition uniqueid.h:53
Declares the Cloud class for storing data and transfering them between worlds.
char * p(char *buf, const char *name, int k, int initial_level, double thresh, int order)
Definition derivatives.cc:72
const std::size_t bufsize
Definition derivatives.cc:16
static bool debug
Definition dirac-hatom.cc:16
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:2643
static const double v
Definition hatom_sf_dirac.cc:20
#define MADNESS_EXCEPTION(msg, value)
Macro for throwing a MADNESS exception.
Definition madness_exception.h:119
#define MADNESS_ASSERT(condition)
Assert a condition that should be free of side-effects since in release builds this might be a no-op.
Definition madness_exception.h:134
#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
Namespace for all elements and tools of MADNESS.
Definition DFParameters.h:10
std::vector< ScalarResult< T > > scalar_result_vector(World &world, std::size_t n)
helper function to create a vector of ScalarResultImpl, circumventing problems with the constructors
Definition macrotaskq.h:163
std::ostream & operator<<(std::ostream &os, const particle< PDIM > &p)
Definition lowrankfunction.h:401
double get_rss_usage_in_GB()
Definition ranks_and_hosts.cpp:10
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:81
@ NodeReplicated
even if there are several ranks per node
Definition worlddc.h:84
@ Distributed
no replication of the container, the container is distributed over the world
Definition worlddc.h:82
@ RankReplicated
replicate the container over all world ranks
Definition worlddc.h:83
void set_impl(std::vector< Function< T, NDIM > > &v, const std::vector< std::shared_ptr< FunctionImpl< T, NDIM > > > vimpl)
Definition vmra.h:738
std::vector< std::shared_ptr< FunctionImpl< T, NDIM > > > get_impl(const std::vector< Function< T, NDIM > > &v)
Definition vmra.h:731
TreeState
Definition funcdefaults.h:59
@ reconstructed
s coeffs at the leaves only
Definition funcdefaults.h:60
@ compressed
d coeffs in internal nodes, s and d coeffs at the root, empty leaves may be present
Definition funcdefaults.h:61
constexpr bool check_tuple_is_valid_task_result()
given a tuple check recursively if all elements are valid task results
Definition macrotaskq.h:222
static const Slice _(0,-1, 1)
static void binary_tuple_loop(tupleT &tuple1, tupleR &tuple2, opT &op)
loop over the tuple elements of both tuples and execute the operation op on each element pair
Definition type_traits.h:742
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
static void unary_tuple_loop(tupleT &tuple, opT &op)
loop over a tuple and apply unary operator op to each element
Definition type_traits.h:732
@ TT_FULL
Definition gentensor.h:120
NDIM & f
Definition mra.h:2622
constexpr bool is_valid_task_result_v
check if type is a valid task result: it must be a WorldObject and must implement gaxpy
Definition macrotaskq.h:208
const Function< T, NDIM > & change_tree_state(const Function< T, NDIM > &f, const TreeState finalstate, bool fence=true)
change tree state of a function
Definition mra.h:2948
decltype(decay_types(std::declval< T >())) decay_tuple
Definition macrotaskpartitioner.h:22
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
constexpr Vector< T, sizeof...(Ts)+1 > vec(T t, Ts... ts)
Factory function for creating a madness::Vector.
Definition vector.h:750
std::string name(const FuncType &type, const int ex=-1)
Definition ccpairfunction.h:28
Function< T, NDIM > copy(const Function< T, NDIM > &f, const std::shared_ptr< WorldDCPmapInterface< Key< NDIM > > > &pmap, bool fence=true)
Create a new copy of the function with different distribution and optional fence.
Definition mra.h:2187
void gaxpy(const double a, ScalarResult< T > &left, const double b, const T &right, const bool fence=true)
the result type of a macrotask must implement gaxpy
Definition macrotaskq.h:244
static const double b
Definition nonlinschro.cc:119
static const double a
Definition nonlinschro.cc:118
static const double L
Definition rk.cc:46
Definition macrotaskq.h:280
friend std::string to_string(const MacroTaskInfo::StoragePolicy sp)
Definition macrotaskq.h:438
StoragePolicy
Definition macrotaskq.h:281
@ StoreFunction
store a madness function in the cloud – can have a large memory impact
Definition macrotaskq.h:282
@ StorePointerToFunction
Definition macrotaskq.h:283
@ StoreFunctionViaPointer
coefficients to the subworlds when the task is started. This is the default policy.
Definition macrotaskq.h:286
static std::vector< MacroTaskInfo > get_all_presets()
helper function to return all presets
Definition macrotaskq.h:350
StoragePolicy storage_policy
Definition macrotaskq.h:426
DistributionType ptr_target_distribution_policy
Definition macrotaskq.h:428
static MacroTaskInfo preset(const std::string name)
Definition macrotaskq.h:313
static std::vector< std::string > get_all_preset_names()
Definition macrotaskq.h:345
static Cloud::StoragePolicy to_cloud_storage_policy(MacroTaskInfo::StoragePolicy policy)
given the MacroTask's storage policy return the corresponding Cloud storage policy
Definition macrotaskq.h:298
DistributionType cloud_distribution_policy
Definition macrotaskq.h:427
friend std::ostream & operator<<(std::ostream &os, const MacroTaskInfo policy)
Definition macrotaskq.h:430
bool check_consistency() const
make sure the policies are consistent
Definition macrotaskq.h:359
friend std::ostream & operator<<(std::ostream &os, const StoragePolicy sp)
Definition macrotaskq.h:290
static StoragePolicy policy_to_string(const std::string policy)
Definition macrotaskq.h:444
nlohmann::json to_json() const
Definition macrotaskq.h:454
void from_vector_of_strings(const std::vector< std::string > &vec)
set policy from a vector of strings, assuming the order is storage policy, cloud distribution policy,...
Definition macrotaskq.h:380
Definition macrotaskq.h:995
World & world
Memoized reference to the world to which this object belongs.
Definition world_object.h:348
World & get_world() const
Definition world_object.h:446
Default load of an object via serialize(ar, t).
Definition archive.h:667
Default store of an object via serialize(ar, t).
Definition archive.h:612
static std::string tolower(std::string s)
make lower case
Definition commandlineparser.h:128
class to temporarily redirect output to cout
Definition print.h:300
RAII class to redirect cout to a file or to /dev/null.
Definition print.h:252
Definition type_traits.h:756
Definition macrotaskq.h:193
Definition macrotaskq.h:178
Definition macrotaskq.h:172
Definition macrotaskq.h:199
Definition macrotaskq.h:187
Definition macrotaskq.h:217
static void load(const Archive &ar, std::shared_ptr< ScalarResultImpl< T > > &ptr)
Definition macrotaskq.h:260
static void store(const Archive &ar, const std::shared_ptr< ScalarResultImpl< T > > &ptr)
Definition macrotaskq.h:250
Definition timing_utilities.h:9
static const double_complex I
Definition tdse1d.cc:164
void doit(World &world)
Definition tdse.cc:921
void e()
Definition test_sig.cc:75
Declares the World class for the parallel runtime environment.
int ProcessID
Used to clearly identify process number/rank.
Definition worldtypes.h:43