@@ -61,40 +61,54 @@ void ModuleList::process(ref_ptr<Candidate> candidate) const {
6161 process((Candidate*) candidate);
6262}
6363
64- void ModuleList::run(Candidate* candidate, bool recursive, bool secondariesFirst) {
64+ void ModuleList::run(Candidate* candidate, bool recursive, bool secondariesFirst, bool waitForSecondaries) {
65+
66+ int taskedSecondaries = 0; //Secondaries already submitted as tasks for all threads
67+
6568 // propagate primary candidate until finished
6669 while (candidate->isActive() && (g_cancel_signal_flag == 0)) {
6770 process(candidate);
6871
6972 // propagate all secondaries before next step of primary
7073 if (recursive and secondariesFirst) {
71- for (size_t i = 0; i < candidate->secondaries.size(); i++) {
74+ #pragma omp taskloop num_tasks(2 * omp_get_max_threads()) untied
75+ for (size_t i = taskedSecondaries; i < candidate->secondaries.size(); i++) {
7276 if (g_cancel_signal_flag != 0)
73- break ;
77+ continue ;
7478 run(candidate->secondaries[i], recursive, secondariesFirst);
7579 }
80+
81+ // increase counter of already tasked secondaries
82+ taskedSecondaries += candidate->secondaries.size();
83+
84+ // wait for completion of tasked secondaries or continue with primary
85+ if (waitForSecondaries) {
86+ #pragma omp taskwait
87+ }
7688 }
7789 }
7890
7991 // propagate secondaries after completing primary
8092 if (recursive and not secondariesFirst) {
81- for (size_t i = 0; i < candidate->secondaries.size(); i++) {
93+ #pragma omp taskloop num_tasks(2 * omp_get_max_threads()) untied
94+ for (size_t i = taskedSecondaries; i < candidate->secondaries.size(); i++) {
8295 if (g_cancel_signal_flag != 0)
83- break ;
96+ continue ;
8497 run(candidate->secondaries[i], recursive, secondariesFirst);
8598 }
99+ #pragma omp taskwait
86100 }
87101
88102 // dump candidae and secondaries if interrupted.
89103 if (candidate->isActive() && (g_cancel_signal_flag != 0))
90104 dumpCandidate(candidate);
91105}
92106
93- void ModuleList::run(ref_ptr<Candidate> candidate, bool recursive, bool secondariesFirst) {
94- run((Candidate*) candidate, recursive, secondariesFirst);
107+ void ModuleList::run(ref_ptr<Candidate> candidate, bool recursive, bool secondariesFirst, bool waitForSecondaries ) {
108+ run((Candidate*) candidate, recursive, secondariesFirst, waitForSecondaries );
95109}
96110
97- void ModuleList::run(const candidate_vector_t *candidates, bool recursive, bool secondariesFirst) {
111+ void ModuleList::run(const candidate_vector_t *candidates, bool recursive, bool secondariesFirst, bool waitForSecondaries ) {
98112 size_t count = candidates->size();
99113
100114#if _OPENMP
@@ -122,7 +136,7 @@ void ModuleList::run(const candidate_vector_t *candidates, bool recursive, bool
122136 }
123137
124138 try {
125- run(candidates->operator[](i), recursive);
139+ run(candidates->operator[](i), recursive, secondariesFirst, waitForSecondaries );
126140 } catch (std::exception &e) {
127141 std::cerr << "Exception in crpropa::ModuleList::run: " << std::endl;
128142 std::cerr << e.what() << std::endl;
@@ -150,12 +164,11 @@ void ModuleList::run(const candidate_vector_t *candidates, bool recursive, bool
150164 }
151165}
152166
153- void ModuleList::run(SourceInterface *source, size_t count, bool recursive, bool secondariesFirst) {
167+ void ModuleList::run(SourceInterface *source, size_t count, bool recursive, bool secondariesFirst, bool waitForSecondaries ) {
154168
155169#if _OPENMP
156170 std::cout << "crpropa::ModuleList: Number of Threads: " << omp_get_max_threads() << std::endl;
157171#endif
158-
159172 ProgressBar progressbar(count);
160173
161174 if (showProgress) {
@@ -189,7 +202,7 @@ void ModuleList::run(SourceInterface *source, size_t count, bool recursive, bool
189202
190203 if (candidate.valid()) {
191204 try {
192- run(candidate, recursive);
205+ run(candidate, recursive, secondariesFirst, waitForSecondaries );
193206 } catch (std::exception &e) {
194207 std::cerr << "Exception in crpropa::ModuleList::run: " << std::endl;
195208 std::cerr << e.what() << std::endl;
0 commit comments