main.cpp 4.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188
  1. #include <unistd.h>
  2. #include <iostream>
  3. #include <err.h>
  4. #include <pthread.h>
  5. #include <poll.h>
  6. #include <pthread.h>
  7. #include <sys/time.h>
  8. #include "messages.h"
  9. #include <signal.h>
  10. #include <map>
  11. #include <limits.h>
  12. #include <memory>
  13. double timediff(
  14. struct timeval const &start_time,
  15. struct timeval const &end_time
  16. )
  17. {
  18. const long long micromill (1000*1000);
  19. const double dmicromill (static_cast<double>(micromill));
  20. struct timeval tdiff;
  21. timersub(&end_time,&start_time,&tdiff);
  22. return static_cast<double>(tdiff.tv_sec)+(static_cast<double>(tdiff.tv_usec)/dmicromill);
  23. }
  24. void* firstThread(void* outlet){
  25. if(outlet){
  26. membuf of;
  27. int fd(*static_cast<int*>(outlet));
  28. struct timeval tv,tve;
  29. gettimeofday(&tv,nullptr);
  30. for(size_t c(0);c<1000;++c){
  31. for(size_t cc(0);cc<1000;++cc){
  32. testMessage t(double(c)+double(cc)/1000);
  33. t.serialize(of);
  34. write(fd,of.buffer(),of.size());
  35. of.reset();
  36. }
  37. }
  38. lastMessage eot;
  39. eot.serialize(of);
  40. write(fd,of.buffer(),of.size());
  41. gettimeofday(&tve,nullptr);
  42. double *tdiff(new double(timediff(tv,tve)));
  43. pthread_exit(tdiff);
  44. }
  45. pthread_exit(nullptr);
  46. }
  47. void handleFds(struct pollfd fd,size_t &ct,size_t& c,bool& eot){
  48. static char buf2[PIPE_BUF];
  49. static membuf ii;
  50. if(fd.revents & POLLIN) {
  51. ++ct;
  52. ii.reset();
  53. size_t n(read(fd.fd,buf2,sizeof(uint32_t)));
  54. if(n==sizeof(uint32_t)){
  55. uint32_t l = (uint32_t)buf2[0] << 24 |
  56. (uint32_t)buf2[1] << 16 |
  57. (uint32_t)buf2[2] << 8 |
  58. (uint32_t)buf2[3];
  59. if(l){
  60. l-=n;
  61. n=read(fd.fd,buf2,l);
  62. if(n){
  63. ii.set(buf2,buf2+n);
  64. std::unique_ptr<pipeMessage> pp(pipeMessage::deserialize(ii));
  65. if(pp.get()){
  66. ++c;
  67. switch(pp->id()){
  68. case 2:
  69. eot= true;
  70. break;
  71. case 1:
  72. if(testMessage* ppp=dynamic_cast<testMessage*>(pp.get())){
  73. double d=ppp->data();
  74. }else {
  75. std::cerr << "unable to cast message #" << pp->id() << std::endl;
  76. }
  77. break;
  78. default:
  79. std::cerr << "unknown message #" << pp->id() << std::endl;
  80. }
  81. }else
  82. std::cerr << "message cast error." << std::endl;
  83. }
  84. else
  85. std::cerr << "error while getting data" << std::endl;
  86. }
  87. }else
  88. std::cerr << "error while getting length" << std::endl;
  89. }
  90. }
  91. size_t c(0),ct(0);
  92. void progress(int){
  93. std::cout << c << "/" << ct<< " messages handled for now." << std::endl;
  94. }
  95. void lostReader(int){
  96. std::cerr << "pipe reader has ended." << std::endl;
  97. }
  98. int main(int argc, char* argv[]){
  99. int nb_threads(2);
  100. if(argc>1){
  101. nb_threads=atoi(argv[1]);
  102. }
  103. #ifdef HAS_SIGINFO
  104. if(signal(SIGINFO,progress)==SIG_ERR)
  105. perror("SIGINFO not supported.");
  106. #endif
  107. if(signal(SIGPIPE,lostReader)==SIG_ERR)
  108. perror("SIGPIPE not supported.");
  109. pthread_t *pid(new pthread_t[nb_threads]);
  110. struct pollfd* fds(new struct pollfd[nb_threads]);
  111. registerMessage(1,testMessage::create);
  112. registerMessage(2,lastMessage::create);
  113. int *fdp(new int[nb_threads*2]);
  114. for(size_t n(0);n<nb_threads;++n){
  115. if(0>pipe2(&fdp[n*2],0)) {
  116. perror("pipe thread.");
  117. exit(-1);
  118. }
  119. pthread_create(&pid[n] ,nullptr,&firstThread,&fdp[n*2+1]);
  120. fds[n].fd=fdp[n*2];
  121. fds[n].events= POLLIN;
  122. }
  123. struct timeval tv,tve;
  124. gettimeofday(&tv,nullptr);
  125. int quit(0);
  126. for(;quit<nb_threads;){
  127. switch(poll(fds,nb_threads,-1)){
  128. case 0:
  129. std::cout << "timeout" << std::endl;
  130. break;
  131. case -1:
  132. perror("Poll");
  133. exit(-1);
  134. break;
  135. default:
  136. for(size_t n(0);n<nb_threads;++n){
  137. bool eot(false);
  138. handleFds(fds[n],ct,c,eot);
  139. if(eot){
  140. std::cout << "Last Message from " << fds[n].fd << std::endl;
  141. ++quit;
  142. }
  143. }
  144. }
  145. }
  146. for(size_t n(0);n<nb_threads;++n){
  147. double *r_thread;
  148. pthread_join(pid[n],reinterpret_cast<void**>(&r_thread));
  149. if(r_thread){
  150. std::cout << "time to write thread " << n <<":\t" << *r_thread << " sec."<< std::endl;
  151. delete r_thread;
  152. }
  153. close(fdp[n*2]);
  154. close(fdp[n*2+1]);
  155. }
  156. delete[] fdp;
  157. delete[] pid;
  158. delete[] fds;
  159. gettimeofday(&tve,nullptr);
  160. std::cout << "time to read " << c << "/" << ct<< " messages:" << timediff(tv,tve)<< " sec." << std::endl;
  161. }