VM2D 1.14
Vortex methods for 2D flows simulation
Loading...
Searching...
No Matches
Queue.cpp
Go to the documentation of this file.
1/*--------------------------------*- VM2D -*-----------------*---------------*\
2| ## ## ## ## #### ##### | | Version 1.14 |
3| ## ## ### ### ## ## ## ## | VM2D: Vortex Method | 2026/03/06 |
4| ## ## ## # ## ## ## ## | for 2D Flow Simulation *----------------*
5| #### ## ## ## ## ## | Open Source Code |
6| ## ## ## ###### ##### | https://www.github.com/vortexmethods/VM2D |
7| |
8| Copyright (C) 2017-2026 I. Marchevsky, K. Sokol, E. Ryatina, A. Kolganova |
9*-----------------------------------------------------------------------------*
10| File name: Queue.cpp |
11| Info: Source code of VM2D |
12| |
13| This file is part of VM2D. |
14| VM2D is free software: you can redistribute it and/or modify it |
15| under the terms of the GNU General Public License as published by |
16| the Free Software Foundation, either version 3 of the License, or |
17| (at your option) any later version. |
18| |
19| VM2D is distributed in the hope that it will be useful, but WITHOUT |
20| ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or |
21| FITNESS FOR A PARTICULAR PURPOSE. See the GNU General Public License |
22| for more details. |
23| |
24| You should have received a copy of the GNU General Public License |
25| along with VM2D. If not, see <http://www.gnu.org/licenses/>. |
26\*---------------------------------------------------------------------------*/
27
28
40#if defined(_WIN32)
41 #include <direct.h>
42#endif
43
44#include <fstream>
45#include <sstream>
46#include <sys/stat.h>
47#include <sys/types.h>
48
49#include "Queue.h"
50
51#include "Preprocessor.h"
52#include "StreamParser.h"
53#include "WorldGen.h"
54
55#ifdef CODE2D
56 #include "Airfoil2D.h"
57 #include "Boundary2D.h"
58 #include "MeasureVP2D.h"
59 #include "Mechanics2D.h"
60 #include "Velocity2D.h"
61 #include "Wake2D.h"
62 #include "WakeDataBase2D.h"
63 #include "World2D.h"
64 #include "Gmres2D.h"
65#endif
66
67#ifdef CODE3D
68 #include "Body3D.h"
69 #include "Passport3D.h"
70 #include "Velocity3D.h"
71 #include "Wake3D.h"
72 #include "World3D.h"
73#endif
74
75using namespace VMlib;
76
77//Конструктор
78Queue::Queue(int& argc, char**& argv)
79{
80#ifdef USE_MPI
81 MPI_Init(&argc, &argv);
82 MPI_Comm_size(MPI_COMM_WORLD, &nProcAll);
83 MPI_Comm_rank(MPI_COMM_WORLD, &myidAll);
84 MPI_Comm_group(MPI_COMM_WORLD, &groupAll);
85#else
86 nProcAll = 1;
87 myidAll = 0;
88 groupAll = -1;
89#endif
90
91 //готовим систему к запуску
92 //(выполняется на глобальном нулевом процессоре)
93 if (myidAll == 0)
94 {
96
98
99 //Устанавливаем флаг занятости в состояние "свободен" всем процессорам,
100#ifdef USE_MPI
101 procState.resize(nProcAll, MPI_UNDEFINED);
102 procStateVar.resize(nProcAll, MPI_UNDEFINED);
103#else
104 procState.resize(nProcAll, -32766);
105 procStateVar.resize(nProcAll, -32766);
106#endif
107
108 numberOfTask.solving = 0; //число решаемых в данный момент задач
109 numberOfTask.prepared = 0; //число подготовленных к запуску задач (в том числе)
110 numberOfTask.finished = 0; //число уже отсчитанных задач
111
112 //текущий номер кванта времени
113 currentKvant = -1;
114 }//if myid==0
115}//Queue()
116
117
118//Деструктор
120{
121 info('i') << "Goodbye!" << std::endl;
122#ifdef USE_MPI
123 MPI_Finalize();
124#endif
125}//~Queue()
126
127
128//Процедура постановки новых задач на отсчет и занятие процессоров
130{
131 //Пересылаем информацию о состоянии процессоров на все компьютеры
132
133#ifdef USE_MPI
134 MPI_Scatter(procState.data(), 1, MPI_INT, &myProcState, 1, MPI_INT, 0, MPI_COMM_WORLD);
135#else
137#endif
138
139#ifdef USE_MPI
140 if (myProcState == MPI_UNDEFINED)
141 world.reset(nullptr);
142#else
143 if (myProcState == -32766)
144 world.reset(nullptr);
145#endif
146
147 if (myidAll == 0)
148 {
149 currentKvant++;
150
151 int nfree = 0; //число свободных процессоров
152
153 //Обнуляем число готовых к старту задач
154 numberOfTask.prepared = 0;
155
156 //считаем число свободных процессоров
157 for (int i = 0; i < nProcAll; ++i)
158#ifdef USE_MPI
159 if (procState[i] == MPI_UNDEFINED)
160#else
161 if (procState[i] == -32766)
162#endif
163 nfree++;
164
165 //формируем номер следующей задачи, которая возможно будет поставлена на отсчет:
166 size_t taskFol = numberOfTask.finished + numberOfTask.solving;
167
168 //проверяем, достаточно ли свободных процессоров для запуска еще одной задачи
169 //при условии, что не все задачи уже решены
170 while ((taskFol < task.size()) && (nfree >= task[taskFol].nProc))
171 {
172 int p = 0; //число найденных свободных процессоров под задачу
173 int j = 0; //текущий счетчик процессоров
174
175 //подбираем свободные процессоры
176 do
177 {
178#ifdef USE_MPI
179 if (procState[j] == MPI_UNDEFINED)
180#else
181 if (procState[j] == -32766)
182#endif
183 {
184 //состояние процессора устанавливаем в "занят текущей задачей"
185 procState[j] = static_cast<int>(taskFol);
186
187 //номер процессора посылаем в перечень процессоров, решающих текущую задачу
188 task[taskFol].proc[p] = j;
189
190 //увеличиваем число найденных процессоров на единицу
191 p++;
192 }//if (ProcState[j]==MPI_UNDEFINED)
193
194 //переходим к следующему процессору
195 j++;
196 } while (p < task[taskFol].nProc);
197
198 //состояние задачи устанавливаем в режим "стартует"
199 task[taskFol].state = TaskState::starting;
200
201 //отмечаем номер кванта, когда задача начала считаться
202 task[taskFol].startEndKvant.first = currentKvant;
203
204 //изменяем счетчики числа задач
205 numberOfTask.prepared++;
206 numberOfTask.solving++;
207
208 //изменяем счетчик свободных процессоров
209 nfree -= task[taskFol].nProc;
210
211 //переходим к следующей задаче
212 taskFol++;
213
214 }//while ((nfree >= Task[task_fol].nproc)&&(task_fol<Task.size()))
215
216 }//if (myidAll == 0)
217
218 //Вывод информации о занятости процессоров задачами
219 if (myidAll == 0)
220 {
221 info.endl();
222 info('i') << "ProcStates: " << std::endl;
223 for (int i = 0; i < nProcAll; ++i)
224 info('-') << "proc[" << i << "] <=> " << ((procState[i] >= 0) ? std::string("problem[" + std::to_string(procState[i]) + "]") : std::string("free")) << std::endl;
225 info.endl();
226 }
227
228 //Пересылаем информацию о состоянии процессоров на все компьютеры
229#ifdef USE_MPI
230 MPI_Scatter(procState.data(), 1, MPI_INT, &myProcState, 1, MPI_INT, 0, MPI_COMM_WORLD);
231#else
233#endif
234
235 //Синхронизируем информацию
236#ifdef USE_MPI
237 MPI_Bcast(&numberOfTask.solving, 1, MPI_INT, 0, MPI_COMM_WORLD);
238 MPI_Bcast(&numberOfTask.prepared, 1, MPI_INT, 0, MPI_COMM_WORLD);
239 MPI_Bcast(&numberOfTask.finished, 1, MPI_INT, 0, MPI_COMM_WORLD);
240#endif
241
242 //список головных процессоров на тех задачах,
243 //которые только что поставлены на счет:
244 std::vector<int> prepList;
245
246 //формируем списки prepList и pspList
247 //(выполняется на глобальном нулевом процессоре)
248 if (myidAll == 0)
249 {
250 for (size_t i = 0; i<task.size(); ++i)
251 //находим задачи в состоянии "стартует"
252 if (task[i].state == TaskState::starting)
253 {
254 //запоминаем номер того процессора, что является там головным
255 prepList.push_back(task[i].proc[0]);
256
257 //состояние задачи вереводим в "считает"
258 task[i].state = TaskState::running;
259 }
260 }//if myidAll==0
261
262
263 //Рассылаем количество стартующих задач в данном кванте на все машины
264#ifdef USE_MPI
265 MPI_Bcast(&numberOfTask.prepared, 1, MPI_INT, 0, MPI_COMM_WORLD);
266#endif
267 if (myidAll > 0)
268 {
269 prepList.resize(numberOfTask.prepared);
270 }//if myid==0
271
272 //следующий фрагмент выполняется только если в данном кванте стартуют новые задачи
273 if (numberOfTask.prepared > 0)
274 {
275 //пересылаем список головных машин стартующих задач на все процессоры
276#ifdef USE_MPI
277 MPI_Bcast(prepList.data(), numberOfTask.prepared, MPI_INT, 0, MPI_COMM_WORLD);
278#endif
279
280 //формируем группу и коммуникатор стартующих головных процессоров
281#ifdef USE_MPI
282 commStarting = MPI_COMM_NULL;
283 MPI_Group_incl(groupAll, numberOfTask.prepared, prepList.data(), &groupStarting);
284 MPI_Comm_create(MPI_COMM_WORLD, groupStarting, &commStarting);
285#else
286 commStarting = 0x04000000;
287 if (numberOfTask.prepared > 0)
288 commStarting = 0;
289#endif
290
291 //следующий фрагмент кода только для стартующих процессов,
292 //которые входят в коммуникатор comm_starting
293#ifdef USE_MPI
294 if (commStarting != MPI_COMM_NULL)
295#else
296 if (commStarting != 0x04000000)
297#endif
298 {
299 //получаем номер процесса в списке стартующих
300 int myidStarting;
301
302#ifdef USE_MPI
303 MPI_Comm_rank(commStarting, &myidStarting);
304#else
305 myidStarting = 0;
306#endif
307
308 //подготовка стартующих задач
309 parallel.myidWork = 0;
310
311#ifdef CODE2D
312 world.reset(new VM2D::World2D(task[myProcState].getPassport()));
313#endif
314
315#ifdef CODE3D
316 world.reset(new VM3D::World3D(task[myProcState].getPassport()));
317#endif
318
319 //Коммуникатор головных процессоров стартующих задач
320 //выполнил свое черное дело и будет удален
321#ifdef USE_MPI
322 MPI_Comm_free(&commStarting);
323#else
324 commStarting = 0x04000000;
325#endif
326 } //if(comm_starting != MPI_COMM_NULL)
327
328 //их группа тоже больше без надобности
329#ifdef USE_MPI
330 MPI_Group_free(&groupStarting);
331#endif
332 }//if (task_prepared>0)
333
334 //формируем группу и коммуникатор считающих головных процессоров
335 //в том числе тех, которые стартуют, и тех, которые ранее входили в commStarting
336 std::vector<int> solvList; //список номеров головных процессов
337
338 if (myidAll == 0)
339 {
340 //Если нулевой процессор свободен
341 //(не является головным в решении задачи в данном кванте) -
342 //- формально присоединяем его в группу головных
343 //необходимо для корректного обмена данными
344#ifdef USE_MPI
345 if (procState[0] == MPI_UNDEFINED)
346#else
347 if (procState[0] == -32766)
348#endif
349 solvList.push_back(0);
350
351 //далее ищем считающие задачи (в том числе только что стартующие)
352 //и собираем номера их головных процессоров
353 //(собираем в порядке просмотра номеров процессоров))
354 for (int s = 0; s<nProcAll; ++s)
355#ifdef USE_MPI
356 if ((procState[s] != MPI_UNDEFINED) && (task[procState[s]].state == TaskState::running) && (task[procState[s]].proc[0] == s))
357#else
358 if ((procState[s] != -32766) && (task[procState[s]].state == TaskState::running) && (task[procState[s]].proc[0] == s))
359#endif
360 solvList.push_back(s);
361
362 sizeCommSolving = static_cast<int>(solvList.size());
363
364 flagFinish.clear();
365 flagFinish.resize(solvList.size());
366
367 }//if myidAll==0
368
369 //число sizeCommSolving на единицу больше, чем taskSolving, если нулевой процессор свободен
370 //и равно taskSolving, если он занят решением задачи
371 //(если он занят решением задачи, то непременно является головным)
372
373 //пересылаем число головных процессоров и их список на все компьютеры
374 //(в том числе головной процессор - независимо от его состояния)
375#ifdef USE_MPI
376 MPI_Bcast(&sizeCommSolving, 1, MPI_INT, 0, MPI_COMM_WORLD);
377#endif
378
379 if (myidAll > 0)
380 solvList.resize(sizeCommSolving);
381
382#ifdef USE_MPI
383 MPI_Bcast(solvList.data(), sizeCommSolving, MPI_INT, 0, MPI_COMM_WORLD);
384#endif
385
386
387#ifdef USE_MPI
388 commSolving = MPI_COMM_NULL;
389
390 //формируем группу и коммуникатор головных процессоров
391 MPI_Group_incl(groupAll, sizeCommSolving, solvList.data(), &groupSolving);
392 MPI_Comm_create(MPI_COMM_WORLD, groupSolving, &commSolving);
393#else
394 commSolving = 0x04000000;
395 if (sizeCommSolving > 0)
396 commSolving = 1;
397#endif
398
399 if (myidAll == 0)
400 {
401 //составление неубывающего массива номеров для верного
402 //распределения задач по процессоров, с тем чтобы задачи
403 //с прошлого кванта оказались в тех же группах, что и раньше
405 }
406
407 //Пересылаем информацию о состоянии процессоров на все компьютеры
408#ifdef USE_MPI
409 MPI_Scatter(procStateVar.data(), 1, MPI_INT, &myProcStateVar, 1, MPI_INT, 0, MPI_COMM_WORLD);
410#else
412#endif
413
414 //Расщепляем весь коммуникатор MPI_COMM_WORLD (т.е. вообще все процессоры)
415 //на коммуникаторы решаемых задач
416 //В результате все процессоры, решающие конкретную задачу,
417 //объединяются в коммуникатор comm_work
418 //Несмотря на то, что он называется везде одинаково, у каждой задачи он свой
419 //и объединяет ровно те процессоры, которые надо
420 //Процессоры, состояние которых установлено в "свободен", т.е.
421 //ProcState=MPI_UNDEFINED (=-32766) ни в один
422 //новый коммуникатор commWork не попадут
423#ifdef USE_MPI
424 MPI_Comm_split(MPI_COMM_WORLD, myProcStateVar, 0, &parallel.commWork);
425#else
426 if (myProcStateVar != -32766)
427 parallel.commWork = 0;
428 else
429 parallel.commWork = 0x04000000;
430#endif
431
432 //следующий код выполняется процессорами, участвующими в решении задач,
433 //следовательно, входящими в коммуникаторы comm_work
434#ifdef USE_MPI
435 if (parallel.commWork != MPI_COMM_NULL)
436#else
437 if (parallel.commWork != 0x04000000)
438#endif
439 {
440 //опеределяем "локальный" номер данного процессора в коммуникаторе commWork и число процессов в нем
441#ifdef USE_MPI
442 MPI_Comm_size(parallel.commWork, &parallel.nProcWork);
443 MPI_Comm_rank(parallel.commWork, &parallel.myidWork);
444#else
446 parallel.myidWork = 0;
447#endif
448 //World равен nullptr только в случае, когда выполняется первый шаг
449 //в этом случае он создан только на главном процессоре группы;
450 //на остальных происходит его создание
451#ifdef CODE2D
452 if (world == nullptr)
453 world.reset(new VM2D::World2D(task[myProcState].getPassport()));
454#endif
455
456#ifdef CODE3D
457 if (world == nullptr)
458 world.reset(new VM3D::World3D(task[myProcState].getPassport()));
459#endif
460
461
462 if ((parallel.myidWork == 0) && (world->getCurrentStep() == 0))
463 {
464#ifdef CODE2D
465 VM2D::World2D& world2D = dynamic_cast<VM2D::World2D&>(*world);
466 //Создание файлов для записи сил
467 for (size_t q = 0; q < world2D.getPassport().airfoilParams.size(); ++q)
468 world2D.GenerateMechanicsHeader(q);
469 //Создание файла для записи временной статистики
470 world2D.getTimers().GenerateStatHeader();
471#endif
472
473#ifdef CODE3D
474 VM3D::World3D& world3D = dynamic_cast<VM3D::World3D&>(*world);
475 //Создание файлов для записи сил
476 //for (size_t q = 0; q < world3D.getPassport().airfoilParams.size(); ++q)
477 // world3D.GenerateMechanicsHeader(q);
478#endif
479 }
480 }
481
482}//TaskSplit()
483
484
485//Процедура обновления состояния задач и процессоров
486void Queue::TaskUpdate()//Обновление состояния задач и процессоров
487{
488 info('i') << "------ Kvant finished ------" << std::endl;
489
490 //Код только для головных процессов, которые входят в коммуникатор commSolving
491 //т.е. выполняется только на головных процессорах, решающих задачи.
492 //Кажется, что можно было бы с тем же эфектом написать здесь if (myidWork==0),
493 //но это не так, поскольку нужно включить в себя еще нулевой процессор,
494 //который в любом случае присоединен к коммуникатору comm_solving,
495 //даже если он не решает задачу (а если решает - то он всегда головной)
496#ifdef USE_MPI
497 if (commSolving != MPI_COMM_NULL)
498 {
499 //определяем номер процессора в данном коммуникаторе -
500 // - вдруг потребуется на будущее!
501 int myidSolving;
502 MPI_Comm_rank(commSolving, &myidSolving);
503
504 //Алгоритм возвращения результатов расчета конкретной задачи
505 //Если срабатывает признак того, что решение задачи можно прекращать
506 //или не решается ни одна задача (а в комуникатор входит
507 //только 0-й процессор, присоединеннй туда насильно)
508 int stopSignal = ((parallel.commWork == MPI_COMM_NULL) || (world->isFinished())) ? 1 : 0;
509
510 //пересылаем признаки прекращения решения задач на нулевой процессор
511 //коммуникатора comm_solving - т.е. на глобальный процессор с номером 0 ---
512 // --- вот для чего его насильно присоединяли к этому коммуникатору
513 //если сам по себе он туда не входил
514 MPI_Gather(&stopSignal, 1, MPI_INT, flagFinish.data(), 1, MPI_INT, 0, commSolving);
515
516 //после пересылки информации коммуникатор comm_solving уничтожается
517 MPI_Comm_free(&commSolving);
518
519 } //if(comm_solving != MPI_COMM_NULL)
520
521 //и соответствующая ему группа тоже удаляется
522 MPI_Group_free(&groupSolving);
523
524
525 if (parallel.commWork != MPI_COMM_NULL)
526 {
527 //К этому месту все процессоры, решающие задачи, приходят уже выполнив все расчеты
528 //в рамках одного кванта времени
529 //Поэтому коммуникаторы решаемых задач нам уже ни к чему - уничтожаем их
530 MPI_Comm_free(&parallel.commWork);
531 }
532#else
533 if (commSolving != 0x04000000)
534 {
535 int myidSolving = 0;
536 int stopSignal = ((parallel.commWork == 0x04000000) || (world->isFinished())) ? 1 : 0;
537 flagFinish[0] = stopSignal;
538 commSolving = 0x04000000;
539 } //if(comm_solving != 0x04000000)
540
541 if (parallel.commWork != 0x04000000)
542 parallel.commWork = 0x04000000;
543#endif
544
545 //Обновление состояния задач и высвобождение процессоров из отсчитавших задач
546 //(выполняется на глобальном нулевом процессоре)
547 if (myidAll == 0)
548 {
549 //Список номеров задач, которые сейчас считаются
550 std::vector<int> taskList;
551
552 //если состояние "считается" запоминаем эту задачу
553 //задачи запоминаются в порядке просмотра процессоров
554 for (int s = 0; s<nProcAll; ++s)
555#ifdef USE_MPI
556 if ( (procState[s] != MPI_UNDEFINED) && (task[procState[s]].state == TaskState::running) && (task[procState[s]].proc[0] == s) )
557#else
558 if ( (procState[s] != -32766) && (task[procState[s]].state == TaskState::running) && (task[procState[s]].proc[0] == s) )
559#endif
560 taskList.push_back(procState[s]);
561
562 //Тем задачам, которые во флаге прекращения счета вернули "1",
563 //ставим состояние завершения счета.
564 //Конструкция +sizeCommSolving-taskSolving введена для смещения в массиве на единицу
565 //если нулевой процессор свободен - тогда он присоединяется формально в нулевом
566 //элементе массива flagFinish, а содержательная часть массива смещается на 1
567 //(в этом случае sizeCommSolving как раз будет на 1 больше, чем taskSolving, см. выше)
568 //если же нулевой процесс занят решением задачи, то смещения нет,
569 //поскольку в этом случае sizeCommSolving и taskSolving равны
570 for (int i = 0; i < numberOfTask.solving; ++i)
571 if (flagFinish[i + sizeCommSolving - numberOfTask.solving] == 1)
572 task[taskList[i]].state = TaskState::finishing;
573 } //if myid==0
574
575 if (myidAll == 0)
576 {
577 for (size_t i = 0; i < task.size(); ++i)
578 {
579 //освобождаем процессоры от сосчитавшихся задач
580 if (task[i].state == TaskState::finishing)
581 {
582 //устанавливаем состояние задачи в "отсчитано"
583 task[i].state = TaskState::done;
584
585 //отмечаем последний квант, когда задача считалась
586 task[i].startEndKvant.second = currentKvant;
587
588 //изменяем счетчики
589 numberOfTask.solving--;
590 numberOfTask.finished++;
591
592 //состояние процессоров, решавших данную задачу, устанавливаем в "свободен"
593 for (int p = 0; p < task[i].nProc; ++p)
594 {
595#ifdef USE_MPI
596 procState[task[i].proc[p]] = MPI_UNDEFINED;
597#else
598 procState[task[i].proc[p]] = -32766;
599#endif
600 }//for p
601 }//if (Task[i].state == 3)
602 }//for i
603
604 //определяем, надо ли еще выделять квант времени или можно заканчивать работу
605 //путем сравнения числа отсчитанных задач с общим числом задач
606 nextKvant = (numberOfTask.finished < static_cast<int>(task.size())) ? 1 : 0;
607 }//if (myid == 0)
608
609#ifdef USE_MPI
610 MPI_Bcast(&nextKvant, 1, MPI_INT, 0, MPI_COMM_WORLD);
611#endif
612}//TaskUpdate()
613
614
615//Процедура, нумерующая задачи в возрастающем порядке
617{
618 std::vector<bool> prFlag;
619 int number = -1;
620
621 prFlag.resize(nProcAll, true); //их надо просматривать
622
623 for (int i = 0; i<nProcAll; ++i)
624
625#ifdef USE_MPI
626 if (procState[i] == MPI_UNDEFINED)
627 {
628 prFlag[i] = false; //уже просмотрели
629 procStateVar[i] = MPI_UNDEFINED;
630 } //if (ProcState[i]==MPI_UNDEFINED)
631#else
632 if (procState[i] == -32766)
633 {
634 prFlag[i] = false; //уже просмотрели
635 procStateVar[i] = -32766;
636 } //if (ProcState[i]==-32766)
637#endif
638
639 for (int i = 0; i<nProcAll; ++i)
640 if (prFlag[i])
641 {
642 prFlag[i] = false;
643 number++;
644
645 for (int s = i; s<nProcAll; ++s)
646 if (procState[s] == procState[i])
647 {
648 procStateVar[s] = number;
649 prFlag[s] = false;
650 }//if (ProcState[s]==ProcState[i])
651 } //if (pr_flag[i])
652}//ConstructProcStateVar()
653
654
655// Запуск вычислительного конвейера (в рамках кванта времени)
657{
658#ifdef USE_MPI
659 if (parallel.commWork != MPI_COMM_NULL)
660#else
661 if (parallel.commWork != 0x04000000)
662#endif
663 {
664 double kvantStartWallTime; //Время на главном процессоре, когда начался очередной квант
665 double deltaWallTime;
666
667 if (parallel.myidWork == 0)
668#ifdef USE_MPI
669 kvantStartWallTime = MPI_Wtime();
670#else
671 kvantStartWallTime = omp_get_wtime();
672#endif
673
675 //if (world->getCurrentStep() == 0)
676 // world->ZeroStep();
677
678 do
679 {
680 if (!world->isFinished())
681 world->Step();
682
683 if (parallel.myidWork == 0)
684#ifdef USE_MPI
685 deltaWallTime = MPI_Wtime() - kvantStartWallTime;
686#else
687 deltaWallTime = omp_get_wtime() - kvantStartWallTime;
688#endif
689
690#ifdef USE_MPI
691 MPI_Bcast(&deltaWallTime, 1, MPI_DOUBLE, 0, parallel.commWork);
692#endif
693 }
694 //Проверка окончания кванта по времени (или завершения задачи)
695 while (deltaWallTime < kvantTime);
696
697 }//if (commWork != MPI_COMM_NULL)
698}//RunConveyer()
699
700
701//Добавление задачи в список
702void Queue::AddTask(int _nProc, std::unique_ptr<PassportGen> _passport)
703{
704 task.resize(task.size() + 1);
705 (task.end() - 1)->nProc = _nProc;
706 (task.end() - 1)->passport = std::move(_passport);
707 (task.end() - 1)->proc.resize(_nProc);
708 (task.end() - 1)->state = TaskState::waiting;
709}//AddTask(...)
710
711
712//Загрузка списка задач
713void Queue::LoadTasksList(const std::string& _tasksFile, const std::string& _mechanicsFile, const std::string& _defaultsFile, const std::string& _switchersFile)
714{
715 std::string extTasksFile = _tasksFile;
716 std::string extMechanicsFile = _mechanicsFile;
717 std::string extDefaultsFile = _defaultsFile;
718 std::string extSwitchersFile = _switchersFile;
719
720 if (
721 fileExistTest(extTasksFile, info, true, { "txt", "TXT" }) &&
722 fileExistTest(extDefaultsFile, info, true, { "txt", "TXT" }) &&
723 fileExistTest(extSwitchersFile, info, true, { "txt", "TXT" })
724 )
725 {
726 std::stringstream tasksFile(Preprocessor(extTasksFile).resultString);
727 std::stringstream defaultsFile(Preprocessor(extDefaultsFile).resultString);
728 std::stringstream switchersFile(Preprocessor(extSwitchersFile).resultString);
729 std::vector<std::string> taskFolder;
730
731 std::unique_ptr<StreamParser> parserTaskList;
732 parserTaskList.reset(new StreamParser(info, "parser", tasksFile, '(', ')'));
733
734 std::vector<std::string> alltasks;
735 parserTaskList->get("problems", alltasks);
736
737 std::ptrdiff_t nTasks = std::count_if(alltasks.begin(), alltasks.end(), [](const std::string& a) {return (a.length() > 0);});
738 info.endl();
739 info('i') << "Number of problems to be solved: " << nTasks << std::endl;
740
741 for (size_t i = 0; i < alltasks.size(); ++i)
742 {
743 if (alltasks[i].length() > 0)
744 {
745 //делим имя задачи + выражение в скобках на 2 подстроки
746 std::pair<std::string, std::string> taskLine = StreamParser::SplitString(info, alltasks[i], false);
747
748 std::string dir = taskLine.first;
749
750 info.endl();
751 info('i') << "-------- Loading problem #" << i << " (" << dir << ") --------" << std::endl;
752
753 //вторую подстроку разделяем на вектор из строк по запятым, стоящим вне фигурных скобок
754 std::vector<std::string> vecTaskLineSecond = StreamParser::StringToVector(taskLine.second, '{', '}');
755
756 //создаем парсер и связываем его с параметрами профиля
757 std::stringstream tskStream(StreamParser::VectorStringToString(vecTaskLineSecond));
758 std::unique_ptr<StreamParser> parserTask;
759
760 parserTask.reset(new StreamParser(info, "problem parameters", tskStream, defaultsFile, {"pspfile", "np", "copyPath"}));
761
762 std::string pspFile;
763 int np;
764
765 //считываем нужные параметры с учетом default-значений
766 parserTask->get("pspfile", pspFile, &defaults::defaultPspFile);
767 parserTask->get("np", np, &defaults::defaultNp);
768
769 if (np > 1) // 31/12/2023
770 {
771 info.endl();
772 info('e') << "Internal MPI-parallelization is not supported longer!" << std::endl;
773 exit(-1);
774 }
775
776 std::string copyPath;
777 parserTask->get("copyPath", copyPath, &defaults::defaultCopyPath);
778 if (copyPath.length() > 0)
779 {
780 //Копировать паспорт поручаем только одному процессу
781 if (myidAll == 0)
782 {
783 std::string initDir(copyPath.begin(), copyPath.end());
784 std::string command;
785
786#if defined(_WIN32)
787 std::replace(dir.begin(), dir.end(), '/', '\\');
788 std::replace(initDir.begin(), initDir.end(), '/', '\\');
789#endif
790
791 VMlib::CreateUserDirectory(dir.c_str(), "");
792
793#if defined(_WIN32)
794 command = "copy \"" + initDir + "\\*.*\" "+ dir + "\\";
795#else
796 command = "cp ./" + initDir + "/* " + dir + "/";
797#endif
798
799 std::cout << "Copying files from folder \"" << initDir << "\" to \"" << dir << "\"" << std::endl;
800
801 int systemRet = system(command.c_str());
802 if(systemRet == -1)
803 {
804 // The system method failed
805 info('e') << "problem #" << i << " (" << dir << \
806 ") copying passport system method failed" << std::endl;
807 exit(-1);
808 }
809 std::cout << "Copying OK " << std::endl << std::endl;
810
811 //copyFile(pspFile, "./" + dir + "/" + pspFile);
812 }
813 }
814
815#ifdef USE_MPI
816 MPI_Barrier(MPI_COMM_WORLD);
817#endif
818
819 std::unique_ptr<PassportGen> ptrPsp;
820
821#ifdef CODE2D
822 ptrPsp.reset(new VM2D::Passport(info, dir, i, "./" + dir + "/" + pspFile, extMechanicsFile, extDefaultsFile, extSwitchersFile, vecTaskLineSecond, {""}));
823#endif
824
825#ifdef CODE3D
826 ptrPsp.reset(new VM3D::Passport(info, dir, i, pspFile, extMechanicsFile, extDefaultsFile, extSwitchersFile, vecTaskLineSecond));
827#endif
828
829 AddTask(np, std::move(ptrPsp));
830
831 info('i') << "-------- Problem #" << i << " (" << dir << ") is loaded --------" << std::endl << std::endl;
832 }
833 } //for i
834
835 //tasksFile.close();
836 tasksFile.clear();
837
838 //defaultsFile.close();
839 defaultsFile.clear();
840
841 //switchersFile.close();
842 switchersFile.clear();
843
844 }
845 else
846 {
847 exit(1);
848 }
849
850}//Queue::LoadTasks()
Заголовочный файл с описанием класса Airfoil.
Заголовочный файл с описанием класса Boundary.
Заголовочный файл с функциями для метода GMRES.
Заголовочный файл с описанием класса MeasureVP.
Заголовочный файл с описанием класса Mechanics.
Заголовочный файл с описанием класса Preprocessor.
Заголовочный файл с описанием класса Queue.
Заголовочный файл с описанием класса StreamParser.
Заголовочный файл с описанием класса Velocity.
Заголовочный файл с описанием класса Wake.
Заголовочный файл с описанием класса WakeDataBase.
Заголовочный файл с описанием класса World2D.
Заголовочный файл с описанием класса WorldGen.
Класс, опеделяющий паспорт двумерной задачи
Definition Passport2D.h:253
Класс, опеделяющий текущую решаемую задачу
Definition World2D.h:77
VMlib::TimersGen & getTimers() const
Возврат ссылки на временную статистику выполнения шага расчета по времени
Definition World2D.h:288
const Passport & getPassport() const
Возврат константной ссылки на паспорт
Definition World2D.h:263
void GenerateMechanicsHeader(size_t mechanicsNumber)
Definition World2D.cpp:2038
void endl()
Вывод в поток логов пустой строки
Definition LogStream.h:103
void assignStream(std::ostream *pStr_, const std::string &label_)
Связывание потока логов с потоком вывода
Definition LogStream.h:80
int nProcWork
Число процессоров, решающих конкретную задачу
Definition Parallel.h:100
int myidWork
Локальный номер процессора, решающего конкретную задачу
Definition Parallel.h:97
int commWork
Коммуникатор для решения конкретной задачи
Definition Parallel.h:93
Класс, позволяющий выполнять предварительную обработку файлов
~Queue()
Деструктор
Definition Queue.cpp:119
void ConstructProcStateVar()
Процедура, нумерующая задачи в возрастающем порядке
Definition Queue.cpp:616
const double kvantTime
Продолжительность кванта времени в секундах
Definition Queue.h:210
Queue(int &argc, char **&argv)
Конструктор
Definition Queue.cpp:78
int sizeCommSolving
Число процессоров в группе для головных процессоров в решаемых в данном кванте времени задачах
Definition Queue.h:166
int myProcState
Состояние данного процессора
Definition Queue.h:196
void LoadTasksList(const std::string &_tasksFile, const std::string &_mechanicsFile, const std::string &_defaultsFile, const std::string &_switchersFile)
Загрузка списка задач
Definition Queue.cpp:713
int myProcStateVar
Состояние данного процессора
Definition Queue.h:203
std::unique_ptr< WorldGen > world
Умный указатель на текущую решаемую задачу
Definition Queue.h:189
int groupSolving
Definition Queue.h:159
int nProcAll
Общее число процессоров
Definition Queue.h:134
std::vector< int > procStateVar
Модифицированный список состояний процессоров
Definition Queue.h:120
struct VMlib::Queue::@0 numberOfTask
Структура, содержащая информацию о количестве задач в данный момент времени
int commSolving
Definition Queue.h:160
std::vector< int > flagFinish
Список возвращаемых флагов останова задачи
Definition Queue.h:128
std::vector< Task > task
Список описаний решаемых задач
Definition Queue.h:186
void RunConveyer()
Запуск вычислительного конвейера (в рамках кванта времени)
Definition Queue.cpp:656
Parallel parallel
Класс, опеделяющий параметры исполнения задачи в параллельном MPI-режиме
Definition Queue.h:228
int currentKvant
Номер текущего кванта времени
Definition Queue.h:213
void AddTask(int _nProc, std::unique_ptr< PassportGen > _passport)
Добавление задачи в список
Definition Queue.cpp:702
void TaskUpdate()
Процедура обновления состояния задач и процессоров
Definition Queue.cpp:486
int myidAll
Глобальный номер процессора
Definition Queue.h:131
int groupAll
Definition Queue.h:156
int groupStarting
Definition Queue.h:157
LogStream info
Поток для вывода логов и сообщений об ошибках
Definition Queue.h:180
int commStarting
Definition Queue.h:158
void TaskSplit()
Процедура постановка новых задач на отсчет и занятие процессоров
Definition Queue.cpp:129
std::vector< int > procState
Список состояний процессоров
Definition Queue.h:108
int nextKvant
Признак необходимости выполнения следующего кванта и продолжения расчета
Definition Queue.h:220
Класс, позволяющий выполнять разбор файлов и строк с настройками и параметрами
static std::vector< std::string > StringToVector(std::string line, char openBracket='(', char closeBracket=')')
Pазбор строки, содержащей запятые, на отдельные строки
static std::pair< std::string, std::string > SplitString(LogStream &info, std::string line, bool upcase=true)
Разбор строки на пару ключ-значение
static std::string VectorStringToString(const std::vector< std::string > &_vecString)
Объединение вектора (списка) из строк в одну строку
void GenerateStatHeader()
Формирование заголовка файла временной статистики
Definition TimesGen.cpp:87
void CreateUserDirectory(const std::string &dir, const std::string &name)
Создание каталога
Definition defs.h:439
void PrintUniversalLogoToStream(std::ostream &str)
Передача в поток вывода универсальной шапки программы VM2D/VM3D.
Definition defs.cpp:92
bool fileExistTest(std::string &fileName, LogStream &info, bool exitKey=false, const std::list< std::string > &extList={})
Проверка существования файла
Definition defs.h:340
@ finishing
задача финиширует
@ starting
задача стартует
@ done
задача решена
@ running
задача решается
@ waiting
задача ожидает запуска
static std::ostream * defaultQueueLogStream
Поток вывода логов и ошибок очереди
Definition defs.h:219
const int defaultNp
Необходимое число процессоров для решения задачи
Definition defs.h:213
const std::string defaultCopyPath
Путь к каталогу с задачей для копирования в новые каталоги
Definition defs.h:216
static std::ostream * defaultWorld2DLogStream
Поток вывода логов и ошибок задачи
Definition defs.h:222
const std::string defaultPspFile
Имя файла с паспортом задачи
Definition defs.h:210