UG

Ubiquity Generator framework

paraCommCPP11.cpp
Go to the documentation of this file.
1/* * * * * * * * * * * * * * * * * * * * * * * * * * * * * * * * * * * * * * */
2/* */
3/* This file is part of the program and software framework */
4/* UG --- Ubquity Generator Framework */
5/* */
6/* Copyright Written by Yuji Shinano <shinano@zib.de>, */
7/* Copyright (C) 2021-2026 by Zuse Institute Berlin, */
8/* licensed under LGPL version 3 or later. */
9/* Commercial licenses are available through <licenses@zib.de> */
10/* */
11/* This code is free software; you can redistribute it and/or */
12/* modify it under the terms of the GNU Lesser General Public License */
13/* as published by the Free Software Foundation; either version 3 */
14/* of the License, or (at your option) any later version. */
15/* */
16/* This program is distributed in the hope that it will be useful, */
17/* but WITHOUT ANY WARRANTY; without even the implied warranty of */
18/* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the */
19/* GNU Lesser General Public License for more details. */
20/* */
21/* You should have received a copy of the GNU Lesser General Public License */
22/* along with this program. If not, see <http://www.gnu.org/licenses/>. */
23/* */
24/* * * * * * * * * * * * * * * * * * * * * * * * * * * * * * * * * * * * * * */
25
26/**@file paraCommCPP11.cpp
27 * @brief ParaComm extension for C++11 thread communication
28 * @author Yuji Shinano
29 *
30 *
31 *
32 */
33
34/*---+----1----+----2----+----3----+----4----+----5----+----6----+----7----+----8----+----9----+----0----+----1----+----2*/
35
36
37#include <cstring>
38#include <cstdlib>
39#ifndef _MSC_VER
40#include <unistd.h>
41#else
42#include <windows.h>
43#endif
44#include "paraCommCPP11.h"
45#include "paraTask.h"
46#include "paraSolution.h"
49#include "paraSolverState.h"
51#include "paraInitialStat.h"
52
53using namespace UG;
54
55std::mutex rankLockMutex;
56
58ParaCommCPP11::threadsTable[ThreadTableSize];
59
60thread_local int
61ParaCommCPP11::localRank = -1; /*< local thread rank */
62
63const char *
64ParaCommCPP11::tagStringTable[] = {
77 TAG_STR(TagRacingRampUpParamSets),
83};
84
85void
86ParaCommCPP11::init( int argc, char **argv )
87{
88 // don't have to take any lock, because only LoadCoordinator call this function
89
90 timer.start();
91 comSize = 0;
92
93 for( int i = 1; i < argc; i++ )
94 {
95 if( strcmp(argv[i], "-sth") == 0 )
96 {
97 i++;
98 if( i < argc )
99 comSize = atoi(const_cast<const char*>(argv[i])); // if -sth 0, then it is considered as use the number of cores system has
100 else
101 {
102 std::cerr << "missing the number of solver threads after parameter '-sth'" << std::endl;
103 exit(1);
104 }
105 }
106 }
107
108 if( comSize > 0 )
109 {
110 comSize++;
111 }
112 else
113 {
114
115#ifdef _MSC_VER
116 SYSTEM_INFO sysinfo;
117 GetSystemInfo(&sysinfo);
118 comSize = sysinfo.dwNumberOfProcessors + 1; //includes logical cpu
119#else
120 comSize = sysconf(_SC_NPROCESSORS_CONF) + 1;
121#endif
122 }
123
124 tokenAccessLockMutex = new std::mutex[comSize];
125 token = new int*[comSize];
126 for( int i = 0; i < comSize; i++ )
127 {
128 token[i] = new int[2];
129 token[i][0] = 0;
130 token[i][1] = -1;
131 }
132
133 /** if you add tag, you should add tagStringTale too */
134 // assert( sizeof(tagStringTable)/sizeof(char*) == N_TH_TAGS );
136
137 /** initialize hashtable */
138 for(int i = 0; i < ThreadTableSize; i++ )
139 {
140 threadsTable[i] = 0;
141 }
142
143 messageQueueTable = new MessageQueueTableElement *[comSize + 1]; // +1 for TimeLimitMonitor
144 sentMessage = new bool[comSize + 1];
145 queueLockMutex = new std::mutex[comSize + 1];
146 sentMsg = new std::condition_variable[comSize + 1];
147 for( int i = 0; i < ( comSize + 1 ); i++ )
148 {
150 sentMessage[i] = false;
151 }
152
153}
154
155void
158 )
159{
160 std::lock_guard<std::mutex> lock(rankLockMutex);
161 assert( localRank == -1 );
162 assert( threadsTable[0] == 0 );
163 localRank = 0;
165 tagTraceFlag = paraParamSet->getBoolParamValue(TagTrace);
166}
167
168void
170 int rank,
172 )
173{
174 std::lock_guard<std::mutex> lock(rankLockMutex);
175 assert( localRank == -1 );
176 assert( threadsTable[rank] == 0 );
177 localRank = rank;
179}
180
181void
183 int rank,
185 )
186{
187 std::lock_guard<std::mutex> lock(rankLockMutex);
188 assert( localRank == -1 );
189 assert( threadsTable[rank] != 0 );
190 localRank = rank;
191 // threadsTable[localRank] = new ThreadsTableElement(localRank, paraParamSet);
192}
193
194void
196 int rank
197 )
198{
199 std::lock_guard<std::mutex> lock(rankLockMutex);
200 assert(rank == localRank);
201 if( threadsTable[rank] == 0 )
202 {
203 THROW_LOGICAL_ERROR2("Invalid remove thread. Rank = ", rank);
204 }
205 else
206 {
207 ThreadsTableElement *elem = threadsTable[rank];
208 delete elem;
209 threadsTable[rank] = 0;
210 localRank = -1;
211 }
212}
213
214bool
216 int rank
217 )
218{
219 // int rank = getRank(); // multi-thread solver may change rank here
220 std::lock_guard<std::mutex> lock(tokenAccessLockMutex[rank]);
221 if( token[rank][0] == rank )
222 {
223 return true;
224 }
225 else
226 {
227 int receivedTag;
228 int source;
229 probe(&source, &receivedTag);
230 TAG_TRACE (Probe, From, source, receivedTag);
231 if( source == 0 && receivedTag == TagToken )
232 {
233 receive(token[rank], 2, ParaINT, 0, TagToken);
234 assert( token[rank][0] == rank );
235 return true;
236 }
237 else
238 {
239 return false;
240 }
241 }
242}
243
244void
246 int rank
247 )
248{
249 // int rank = getRank(); // multi-thread solver may change rank here
250 std::lock_guard<std::mutex> lock(tokenAccessLockMutex[rank]);
251 assert( token[rank][0] == rank && rank != 0 );
252 token[rank][0] = ( token[rank][0] % (comSize - 1) ) + 1;
253 token[rank][1] = -1;
254 send(token[rank], 2, ParaINT, 0, TagToken);
255}
256
257bool
259 int rank
260 )
261{
262 // int rank = getRank(); // multi-thread solver may change rank here
263 std::lock_guard<std::mutex> lock(tokenAccessLockMutex[rank]);
264 if( rank == token[rank][0] )
265 {
266 if( token[rank][1] == token[rank][0] ) token[rank][1] = -2;
267 else if( token[rank][1] == -1 ) token[rank][1] = token[rank][0];
268 token[rank][0] = ( token[rank][0] % (comSize - 1) ) + 1;
269 }
270 else
271 {
272 THROW_LOGICAL_ERROR4("Invalid token update. Rank = ", getRank(), ", token = ", token[0] );
273 }
274 send(token[rank], 2, ParaINT, 0, TagToken);
275 if( token[rank][1] == -2 )
276 {
277 return true;
278 }
279 else
280 {
281 return false;
282 }
283}
284
285void
287 int rank,
288 int *inToken
289 )
290{
291 // int rank = getRank();
292 std::lock_guard<std::mutex> lock(tokenAccessLockMutex[rank]);
293 assert( rank == 0 || ( rank != 0 && inToken[0] == rank ) );
294 token[rank][0] = inToken[0];
295 token[rank][1] = inToken[1];
296}
297
298
299
301{
302 std::lock_guard<std::mutex> lock(rankLockMutex);
303 for(int i = 0; i < ThreadTableSize; i++ )
304 {
305 if( threadsTable[i] )
306 {
307 delete threadsTable[i];
308 }
309 }
310
311 for( int i = 0; i < comSize; i++ )
312 {
313 delete [] token[i];
314 }
315 delete [] token;
316 delete [] tokenAccessLockMutex;
317
318 for(int i = 0; i < (comSize + 1); i++)
319 {
321 while( elem )
322 {
323 if( elem->getData() )
324 {
325 if( !freeStandardTypes(elem) )
326 {
327 ABORT_LOGICAL_ERROR2("Requested type is not implemented. Type = ", elem->getDataTypeId() );
328 }
329 }
330 delete elem;
332 }
333 delete messageQueueTable[i];
334 }
335 delete [] messageQueueTable;
336
337 if( sentMessage ) delete [] sentMessage;
338 if( queueLockMutex ) delete [] queueLockMutex;
339 if( sentMsg ) delete [] sentMsg;
340
341}
342
343int
345 )
346{
347 // no lock needed! localRank is thread_local the read needs no sync
348 if( localRank >= 0 ) return localRank;
349 else return -1; // No ug threads
350}
351
352std::ostream *
354 )
355{
356 //no lock needed here!
357 if( !threadsTable[localRank] ) return 0;
358 // assert( threadsTable[localRank] );
360}
361
362void *
364 const void* buffer,
365 int count,
366 const int datatypeId
367 )
368{
369 void *newBuf = 0;
370 if( count == 0 ) return newBuf;
371
372 switch(datatypeId)
373 {
374 case ParaCHAR :
375 {
376 newBuf = new char[count];
377 memcpy(newBuf, buffer, (unsigned long int)sizeof(char)*count);
378 break;
379 }
380 case ParaSHORT :
381 {
382 newBuf = new short[count];
383 memcpy(newBuf, buffer, (unsigned long int)sizeof(short)*count);
384 break;
385 }
386 case ParaINT :
387 {
388 newBuf = new int[count];
389 memcpy(newBuf, buffer, (unsigned long int)sizeof(int)*count);
390 break;
391 }
392 case ParaLONG :
393 {
394 newBuf = new long[count];
395 memcpy(newBuf, buffer, (unsigned long int)sizeof(long)*count);
396 break;
397 }
398 case ParaUNSIGNED_CHAR :
399 {
400 newBuf = new unsigned char[count];
401 memcpy(newBuf, buffer, (unsigned long int)sizeof(unsigned char)*count);
402 break;
403 }
404 case ParaUNSIGNED_SHORT :
405 {
406 newBuf = new unsigned short[count];
407 memcpy(newBuf, buffer, (unsigned long int)sizeof(unsigned short)*count);
408 break;
409 }
410 case ParaUNSIGNED :
411 {
412 newBuf = new unsigned int[count];
413 memcpy(newBuf, buffer, (unsigned long int)sizeof(unsigned int)*count);
414 break;
415 }
416 case ParaUNSIGNED_LONG :
417 {
418 newBuf = new unsigned long[count];
419 memcpy(newBuf, buffer, (unsigned long int)sizeof(unsigned long)*count);
420 break;
421 }
422 case ParaFLOAT :
423 {
424 newBuf = new float[count];
425 memcpy(newBuf, buffer, (unsigned long int)sizeof(float)*count);
426 break;
427 }
428 case ParaDOUBLE :
429 {
430 newBuf = new double[count];
431 memcpy(newBuf, buffer, (unsigned long int)sizeof(double)*count);
432 break;
433 }
434 case ParaLONG_DOUBLE :
435 {
436 newBuf = new long double[count];
437 memcpy(newBuf, buffer, (unsigned long int)sizeof(long double)*count);
438 break;
439 }
440 case ParaBYTE :
441 {
442 newBuf = new char[count];
443 memcpy(newBuf, buffer, (unsigned long int)sizeof(char)*count);
444 break;
445 }
446 case ParaSIGNED_CHAR :
447 {
448 newBuf = new char[count];
449 memcpy(newBuf, buffer, (unsigned long int)sizeof(char)*count);
450 break;
451 }
452 case ParaLONG_LONG :
453 {
454 newBuf = new long long[count];
455 memcpy(newBuf, buffer, (unsigned long int)sizeof(long long)*count);
456 break;
457 }
459 {
460 newBuf = new unsigned long long[count];
461 memcpy(newBuf, buffer, (unsigned long int)sizeof(unsigned long long)*count);
462 break;
463 }
464 case ParaBOOL :
465 {
466 newBuf = new bool[count];
467 memcpy(newBuf, buffer, (unsigned long int)sizeof(bool)*count);
468 break;
469 }
470 default :
471 THROW_LOGICAL_ERROR2("This type is not implemented. Type = ", datatypeId);
472 }
473
474 return newBuf;
475}
476
477void
479 void *dest, const void *src, int count, int datatypeId
480 )
481{
482
483 if( count == 0 ) return;
484
485 switch(datatypeId)
486 {
487 case ParaCHAR :
488 {
489 memcpy(dest, src, (unsigned long int)sizeof(char)*count);
490 break;
491 }
492 case ParaSHORT :
493 {
494 memcpy(dest, src, (unsigned long int)sizeof(short)*count);
495 break;
496 }
497 case ParaINT :
498 {
499 memcpy(dest, src, (unsigned long int)sizeof(int)*count);
500 break;
501 }
502 case ParaLONG :
503 {
504 memcpy(dest, src, (unsigned long int)sizeof(long)*count);
505 break;
506 }
507 case ParaUNSIGNED_CHAR :
508 {
509 memcpy(dest, src, (unsigned long int)sizeof(unsigned char)*count);
510 break;
511 }
512 case ParaUNSIGNED_SHORT :
513 {
514 memcpy(dest, src, (unsigned long int)sizeof(unsigned short)*count);
515 break;
516 }
517 case ParaUNSIGNED :
518 {
519 memcpy(dest, src, (unsigned long int)sizeof(unsigned int)*count);
520 break;
521 }
522 case ParaUNSIGNED_LONG :
523 {
524 memcpy(dest, src, (unsigned long int)sizeof(unsigned long)*count);
525 break;
526 }
527 case ParaFLOAT :
528 {
529 memcpy(dest, src, (unsigned long int)sizeof(float)*count);
530 break;
531 }
532 case ParaDOUBLE :
533 {
534 memcpy(dest, src, (unsigned long int)sizeof(double)*count);
535 break;
536 }
537 case ParaLONG_DOUBLE :
538 {
539 memcpy(dest, src, (unsigned long int)sizeof(long double)*count);
540 break;
541 }
542 case ParaBYTE :
543 {
544 memcpy(dest, src, (unsigned long int)sizeof(char)*count);
545 break;
546 }
547 case ParaSIGNED_CHAR :
548 {
549 memcpy(dest, src, (unsigned long int)sizeof(char)*count);
550 break;
551 }
552 case ParaLONG_LONG :
553 {
554 memcpy(dest, src, (unsigned long int)sizeof(long long)*count);
555 break;
556 }
558 {
559 memcpy(dest, src, (unsigned long int)sizeof(unsigned long long)*count);
560 break;
561 }
562 case ParaBOOL :
563 {
564 memcpy(dest, src, (unsigned long int)sizeof(bool)*count);
565 break;
566 }
567 default :
568 THROW_LOGICAL_ERROR2("This type is not implemented. Type = ", datatypeId);
569 }
570
571}
572
573void
575 void* buffer,
576 int count,
577 const int datatypeId
578 )
579{
580
581 if( count == 0 ) return;
582
583 switch(datatypeId)
584 {
585 case ParaCHAR :
586 {
587 delete [] static_cast<char *>(buffer);
588 break;
589 }
590 case ParaSHORT :
591 {
592 delete [] static_cast<short *>(buffer);
593 break;
594 }
595 case ParaINT :
596 {
597 delete [] static_cast<int *>(buffer);
598 break;
599 }
600 case ParaLONG :
601 {
602 delete [] static_cast<long *>(buffer);
603 break;
604 }
605 case ParaUNSIGNED_CHAR :
606 {
607 delete [] static_cast<unsigned char *>(buffer);
608 break;
609 }
610 case ParaUNSIGNED_SHORT :
611 {
612 delete [] static_cast<unsigned short *>(buffer);
613 break;
614 }
615 case ParaUNSIGNED :
616 {
617 delete [] static_cast<unsigned int *>(buffer);
618 break;
619 }
620 case ParaUNSIGNED_LONG :
621 {
622 delete [] static_cast<unsigned long *>(buffer);
623 break;
624 }
625 case ParaFLOAT :
626 {
627 delete [] static_cast<float *>(buffer);
628 break;
629 }
630 case ParaDOUBLE :
631 {
632 delete [] static_cast<double *>(buffer);
633 break;
634 }
635 case ParaLONG_DOUBLE :
636 {
637 delete [] static_cast<long double *>(buffer);
638 break;
639 }
640 case ParaBYTE :
641 {
642 delete [] static_cast<char *>(buffer);
643 break;
644 }
645 case ParaSIGNED_CHAR :
646 {
647 delete [] static_cast<char *>(buffer);
648 break;
649 }
650 case ParaLONG_LONG :
651 {
652 delete [] static_cast<long long *>(buffer);
653 break;
654 }
656 {
657 delete [] static_cast<unsigned long long *>(buffer);
658 break;
659 }
660 case ParaBOOL :
661 {
662 delete [] static_cast<bool *>(buffer);;
663 break;
664 }
665 default :
666 THROW_LOGICAL_ERROR2("This type is not implemented. Type = ", datatypeId);
667 }
668
669}
670
671bool
673 MessageQueueElement *elem ///< pointer to a message queue element
674 )
675{
676 if( elem->getDataTypeId() < UG_USER_TYPE_FIRST )
677 {
678 freeMem(elem->getData(), elem->getCount(), elem->getDataTypeId() );
679 }
680 else
681 {
682 switch( elem->getDataTypeId())
683 {
684 case ParaInstanceType:
685 {
686 delete reinterpret_cast<ParaInstance *>(elem->getData());
687 break;
688 }
689 case ParaSolutionType:
690 {
691 delete reinterpret_cast<ParaSolution *>(elem->getData());
692 break;
693 }
694 case ParaParamSetType:
695 {
696 delete reinterpret_cast<ParaParamSet *>(elem->getData());
697 break;
698 }
699 case ParaTaskType:
700 {
701 delete reinterpret_cast<ParaTask *>(elem->getData());
702 break;
703 }
705 {
706 delete reinterpret_cast<ParaSolverState *>(elem->getData());
707 break;
708 }
710 {
711 delete reinterpret_cast<ParaCalculationState *>(elem->getData());
712 break;
713 }
715 {
716 delete reinterpret_cast<ParaSolverTerminationState *>(elem->getData());
717 break;
718 }
720 {
721 delete reinterpret_cast<ParaRacingRampUpParamSet *>(elem->getData());
722 break;
723 }
724 default:
725 {
726 return false;
727 }
728 }
729 }
730 return true;
731}
732
733bool
735 )
736{
737 // std::cout << "size = " << sizeof(tagStringTable)/sizeof(char*) << ", N_TH_TAGS = " << N_TH_TAGS << std::endl;
738 return ( sizeof(tagStringTable)/sizeof(char*) == N_TH_TAGS );
739}
740
741const char *
743 int tag /// tag to be converted to string
744 )
745{
746 assert( tag >= 0 && tag < N_TH_TAGS );
747 return tagStringTable[tag];
748}
749
750
751int
753 void* buffer,
754 int count,
755 const int datatypeId,
756 int root
757 )
758{
759 if( getRank() == root )
760 {
761 for(int i=0; i < comSize; i++)
762 {
763 if( i != root )
764 {
765 send(buffer, count, datatypeId, i, -1);
766 }
767 }
768 }
769 else
770 {
771 receive(buffer, count, datatypeId, root, -1);
772 }
773 return 0;
774}
775
776int
778 void* buffer,
779 int count,
780 const int datatypeId,
781 int dest,
782 const int tag
783 )
784{
785 {
786 std::lock_guard<std::mutex> lock(queueLockMutex[dest]);
787 messageQueueTable[dest]->enqueue(sentMsg[dest], queueLockMutex[dest], &sentMessage[dest],
788 new MessageQueueElement(getRank(), count, datatypeId, tag,
789 allocateMemAndCopy(buffer, count, datatypeId) ) );
790 }
791 TAG_TRACE (Send, To, dest, tag);
792 return 0;
793}
794
795int
797 void* buffer,
798 int count,
799 const int datatypeId,
800 int source,
801 const int tag
802 )
803{
804 int qRank = getRank();
805 MessageQueueElement *elem = 0;
806 if( !messageQueueTable[qRank]->checkElement(source, datatypeId, tag) )
807 {
808 messageQueueTable[qRank]->waitMessage(sentMsg[qRank], queueLockMutex[qRank], &sentMessage[qRank], source, datatypeId, tag);
809 }
810 {
811 std::lock_guard<std::mutex> lock(queueLockMutex[qRank]);
812 elem = messageQueueTable[qRank]->extarctElement(&sentMessage[qRank],source, datatypeId, tag);
813 }
814 assert(elem);
815 copy( buffer, elem->getData(), count, datatypeId );
816 freeMem(elem->getData(), count, datatypeId );
817 delete elem;
818 TAG_TRACE (Recv, From, source, tag);
819 return 0;
820}
821
822void
824 const int source,
825 const int tag,
826 int *receivedTag
827 )
828{
829 /*
830 // Just wait, iProbe and receive will be performed after this call
831 messageQueueTable[getRank()]->waitMessage(source, datatypeId, tag);
832 TAG_TRACE (Probe, From, source, tag);
833 return 0;
834 */
835 int qRank = getRank();
836 // LOCKED ( &queueLock[getRank()] )
837 // {
838 // messageQueueTable[qRank]->waitMessage(sentMsg[qRank], queueLockMutex[qRank], &sentMessage[qRank], source, receivedTag);
839 // }
840 (*receivedTag) = tag;
841 messageQueueTable[qRank]->waitMessage(sentMsg[qRank], queueLockMutex[qRank], &sentMessage[qRank], source, receivedTag);
842 TAG_TRACE (Probe, From, source, *receivedTag);
843 return;
844}
845
846bool
848 int* source,
849 int* tag
850 )
851{
852 int qRank = getRank();
853 messageQueueTable[qRank]->waitMessage(sentMsg[qRank], queueLockMutex[qRank], &sentMessage[qRank]);
855 *source = elem->getSource();
856 *tag = elem->getTag();
857 TAG_TRACE (Probe, From, *source, *tag);
858 return true;
859}
860
861bool
863 int* source,
864 int* tag
865 )
866{
867 bool flag = false;
868 int qRank = getRank();
869 {
870 std::lock_guard<std::mutex> lock(queueLockMutex[qRank]);
871 flag = !(messageQueueTable[qRank]->isEmpty());
872 if( flag )
873 {
874 if( *tag == TagAny )
875 {
877 *source = elem->getSource();
878 *tag = elem->getTag();
879 TAG_TRACE (Iprobe, From, *source, *tag);
880 }
881 else
882 {
884 if( elem )
885 {
886 *source = elem->getSource();
887 *tag = elem->getTag();
888 TAG_TRACE (Iprobe, From, *source, *tag);
889 // flag = true;
890 }
891 else
892 {
893 flag = false;
894 }
895 }
896 }
897 }
898 return flag;
899}
900
901int
903 void* buffer,
904 const int datatypeId,
905 int dest,
906 const int tag
907 )
908{
909 {
910 std::lock_guard<std::mutex> lock(queueLockMutex[dest]);
911 messageQueueTable[dest]->enqueue(sentMsg[dest], queueLockMutex[dest], &sentMessage[dest],
912 new MessageQueueElement(getRank(), 1, datatypeId, tag, buffer ) );
913 }
914 TAG_TRACE (Send, To, dest, tag);
915 return 0;
916}
917
918int
920 void** buffer,
921 const int datatypeId,
922 int source,
923 const int tag
924 )
925{
926 int qRank = getRank();
927 if( !messageQueueTable[qRank]->checkElement(source, datatypeId, tag) )
928 {
929 messageQueueTable[qRank]->waitMessage(sentMsg[qRank], queueLockMutex[qRank], &sentMessage[qRank], source, datatypeId, tag);
930 }
931 MessageQueueElement *elem = 0;
932 {
933 std::lock_guard<std::mutex> lock(queueLockMutex[qRank]);
934 elem = messageQueueTable[qRank]->extarctElement(&sentMessage[qRank], source, datatypeId, tag);
935 }
936 assert(elem);
937 *buffer = elem->getData();
938 delete elem;
939 TAG_TRACE (Recv, From, source, tag);
940 return 0;
941}
Class for message queue element.
int getSource()
getter of source rank
int getDataTypeId()
getter of the data type id
void * getData()
getter of data
int getTag()
getter of the message tag
int getCount()
getter of the number of the data type elements
Class of MessageQueueTableElement.
void enqueue(std::condition_variable &sentMsg, std::mutex &queueLockMutex, bool *sentMessage, MessageQueueElement *newElement)
enqueue a message
MessageQueueElement * extarctElement(bool *sentMessage, int source, int datatypeId, int tag)
extracts a message
MessageQueueElement * getHead()
getter of head
void waitMessage(std::condition_variable &sentMsg, std::mutex &queueLockMutex, bool *sentMessage)
wait for a message coming to a queue
MessageQueueElement * checkElementWithTag(int tag)
check if the specified message with tag exists or nor
bool isEmpty()
check if the queue is empty or not
Abstract interface for calculation state in a ParaSolver. Holds no data; derived classes own the fiel...
virtual void solverInit(ParaParamSet *paraParamSet)
initializer for Solvers
int comSize
communicator size : number of threads joined in this system
bool probe(int *source, int *tag)
probe function which waits a new message
int send(void *bufer, int count, const int datatypeId, int dest, const int tag)
send function for standard ParaData types
int ** token
index 0: token index 1: token color -1: green > 0: yellow ( termination origin solver number ) -2: re...
bool iProbe(int *source, int *tag)
iProbe function which checks if a new message is arrived or not
std::mutex rankLockMutex
mutex to access rank
int receive(void *bufer, int count, const int datatypeId, int source, const int tag)
receive function for standard ParaData types
void freeMem(void *buffer, int count, const int datatypeId)
free memory
std::ostream * getOstream()
get ostream pointer
virtual bool tagStringTableIsSetUpCoorectly()
check if tag string table (for debugging) set up correctly
bool tagTraceFlag
indicate if tags are traced or not
void * allocateMemAndCopy(const void *buffer, int count, const int datatypeId)
allocate memory and copy message
int uTypeSend(void *bufer, const int datatypeId, int dest, int tag)
User type send for created data type.
virtual const char * getTagString(int tag)
get Tag string for debugging
ParaSysTimer timer
system timer
int uTypeReceive(void **bufer, const int datatypeId, int source, int tag)
User type receive for created data type.
virtual void passToken(int rank)
pass token to from the rank to the next
MessageQueueTableElement ** messageQueueTable
message queue table
static const char * tagStringTable[]
tag name string table
int getRank()
get rank of caller's thread
virtual void solverDel(int rank)
delete Solver from this communicator
bool freeStandardTypes(MessageQueueElement *elem)
free memory
bool * sentMessage
sent message flag for synchronization
void copy(void *dest, const void *src, int count, int datatypeId)
copy message
virtual ~ParaCommCPP11()
destructor of this communicator
static thread_local int localRank
local thread rank
virtual void solverReInit(int rank, ParaParamSet *paraParamSet)
reinitializer of a specific Solver
virtual bool passTermToken(int rank)
pass termination token from the rank to the next
virtual void setToken(int rank, int *inToken)
set received token to this communicator
static ThreadsTableElement * threadsTable[ThreadTableSize]
threads table: index is thread rank
int bcast(void *buffer, int count, const int datatypeId, int root)
broadcast function for standard ParaData types
void waitSpecTagFromSpecSource(const int source, const int tag, int *receivedTag)
wait function for a specific tag from a specific source coming from
virtual void lcInit(ParaParamSet *paraParamSet)
initializer for LoadCoordinator
std::mutex * queueLockMutex
mutex for synchronization
std::mutex * tokenAccessLockMutex
mutex to access token
std::condition_variable * sentMsg
condition variable for synchronization
virtual bool waitToken(int rank)
wait token when UG runs with deterministic mode
class for instance data
Definition: paraInstance.h:51
class ParaParamSet
Definition: paraParamSet.h:850
class ParaRacingRampUpParamSet (parameter set for racing ramp-up)
class for solution
Definition: paraSolution.h:54
class ParaSolverState (ParaSolver state object for notification message)
class ParaSolverTerminationState (Solver termination state in a ParaSolver)
void start(void)
start timer
class ParaTask
Definition: paraTask.h:542
Class of ThreadsTableElement.
std::ostream * getOstream()
getter of tag trace stream of this rank
static ScipParaParamSet * paraParamSet
Definition: fscip.cpp:74
static const int ParaUNSIGNED_LONG
Definition: paraComm.h:73
static const int TagAckCompletion
Definition: paraTagDef.h:62
static const int TagCompletionOfCalculation
Definition: paraTagDef.h:54
static const int TagWinner
Definition: paraTagDef.h:60
static const int ParaTaskType
static const int ParaInstanceType
Definition: paraCommCPP11.h:98
static const int TagSolution
Definition: paraTagDef.h:51
static const int ParaUNSIGNED_SHORT
Definition: paraComm.h:71
static const int TagToken
Definition: paraTagDef.h:63
static const int TagTaskReceived
Definition: paraTagDef.h:48
static const int N_TH_TAGS
Definition: paraTagDef.h:85
static const int TagInterruptRequest
Definition: paraTagDef.h:57
static const int TagNotificationId
Definition: paraTagDef.h:55
static const int ParaParamSetType
static const int TagIncumbentValue
Definition: paraTagDef.h:52
static const int ParaLONG_DOUBLE
Definition: paraComm.h:77
static const int ParaINT
Definition: paraComm.h:66
static const int TagTerminated
Definition: paraTagDef.h:58
static const int ParaSolverStateType
static const int ParaCalculationStateType
static const int ParaLONG
Definition: paraComm.h:67
static const int TagTerminateRequest
Definition: paraTagDef.h:56
static const int TagAny
Definition: paraTagDef.h:44
static const int ParaBYTE
Definition: paraComm.h:79
static const int ParaSolverTerminationStateType
static const int ParaUNSIGNED
Definition: paraComm.h:72
static const int TagRampUp
Definition: paraTagDef.h:50
static const int TagSolverState
Definition: paraTagDef.h:53
static const int TagHardTimeLimit
Definition: paraTagDef.h:61
static const int ParaFLOAT
Definition: paraComm.h:75
static const int ParaBOOL
Definition: paraComm.h:78
static const int ParaCHAR
Definition: paraComm.h:64
static const int UG_USER_TYPE_FIRST
user defined transfer data types
Definition: paraCommCPP11.h:97
static const int TagDiffSubproblem
Definition: paraTagDef.h:49
static const int ThreadTableSize
size of thread table : this limits the number of threads
static const int TagTask
Definition: paraTagDef.h:47
static const int ParaSolutionType
Definition: paraCommCPP11.h:99
static const int ParaRacingRampUpParamType
static const int TagParaInstance
Definition: paraTagDef.h:82
static const int ParaSHORT
Definition: paraComm.h:65
static const int TagTrace
Definition: paraParamSet.h:72
static const int ParaUNSIGNED_LONG_LONG
Definition: paraComm.h:74
static const int ParaLONG_LONG
Definition: paraComm.h:68
static const int ParaUNSIGNED_CHAR
Definition: paraComm.h:70
static const int ParaDOUBLE
Definition: paraComm.h:76
static const int ParaSIGNED_CHAR
Definition: paraComm.h:69
Base class for calculation state.
std::mutex rankLockMutex
ParaComm extension for C++11 thread communication.
#define TAG_TRACE(call, fromTo, sourceDest, tag)
Definition: paraCommCPP11.h:58
#define ABORT_LOGICAL_ERROR2(msg1, msg2)
Definition: paraDef.h:78
#define THROW_LOGICAL_ERROR2(msg1, msg2)
Definition: paraDef.h:69
#define THROW_LOGICAL_ERROR4(msg1, msg2, msg3, msg4)
Definition: paraDef.h:103
Base class for initial statistics collecting class.
Base class for racing ramp-up parameter set.
Base class for solution.
This class has solver state to be transferred.
This class contains solver termination state which is transferred form Solver to LC.
#define TAG_STR(tag)
Definition: paraTagDef.h:40
Base class for ParaTask.