| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188 |
- #include <unistd.h>
- #include <iostream>
- #include <err.h>
- #include <pthread.h>
- #include <poll.h>
- #include <pthread.h>
- #include <sys/time.h>
- #include "messages.h"
- #include <signal.h>
- #include <map>
- #include <limits.h>
- #include <memory>
- double timediff(
- struct timeval const &start_time,
- struct timeval const &end_time
- )
- {
- const long long micromill (1000*1000);
- const double dmicromill (static_cast<double>(micromill));
- struct timeval tdiff;
- timersub(&end_time,&start_time,&tdiff);
- return static_cast<double>(tdiff.tv_sec)+(static_cast<double>(tdiff.tv_usec)/dmicromill);
- }
- void* firstThread(void* outlet){
- if(outlet){
- membuf of;
- int fd(*static_cast<int*>(outlet));
- struct timeval tv,tve;
- gettimeofday(&tv,nullptr);
- for(size_t c(0);c<1000;++c){
- for(size_t cc(0);cc<1000;++cc){
- testMessage t(double(c)+double(cc)/1000);
- t.serialize(of);
- write(fd,of.buffer(),of.size());
- of.reset();
- }
- }
- lastMessage eot;
- eot.serialize(of);
- write(fd,of.buffer(),of.size());
- gettimeofday(&tve,nullptr);
- double *tdiff(new double(timediff(tv,tve)));
- pthread_exit(tdiff);
- }
- pthread_exit(nullptr);
- }
- void handleFds(struct pollfd fd,size_t &ct,size_t& c,bool& eot){
- static char buf2[PIPE_BUF];
- static membuf ii;
- if(fd.revents & POLLIN) {
- ++ct;
- ii.reset();
- size_t n(read(fd.fd,buf2,sizeof(uint32_t)));
- if(n==sizeof(uint32_t)){
- uint32_t l = (uint32_t)buf2[0] << 24 |
- (uint32_t)buf2[1] << 16 |
- (uint32_t)buf2[2] << 8 |
- (uint32_t)buf2[3];
- if(l){
- l-=n;
- n=read(fd.fd,buf2,l);
- if(n){
- ii.set(buf2,buf2+n);
- std::unique_ptr<pipeMessage> pp(pipeMessage::deserialize(ii));
- if(pp.get()){
- ++c;
- switch(pp->id()){
- case 2:
- eot= true;
- break;
- case 1:
- if(testMessage* ppp=dynamic_cast<testMessage*>(pp.get())){
- double d=ppp->data();
- }else {
- std::cerr << "unable to cast message #" << pp->id() << std::endl;
- }
- break;
- default:
- std::cerr << "unknown message #" << pp->id() << std::endl;
- }
- }else
- std::cerr << "message cast error." << std::endl;
- }
- else
- std::cerr << "error while getting data" << std::endl;
- }
- }else
- std::cerr << "error while getting length" << std::endl;
- }
- }
- size_t c(0),ct(0);
- void progress(int){
- std::cout << c << "/" << ct<< " messages handled for now." << std::endl;
- }
- void lostReader(int){
- std::cerr << "pipe reader has ended." << std::endl;
- }
- int main(int argc, char* argv[]){
- int nb_threads(2);
- if(argc>1){
- nb_threads=atoi(argv[1]);
- }
- #ifdef HAS_SIGINFO
- if(signal(SIGINFO,progress)==SIG_ERR)
- perror("SIGINFO not supported.");
- #endif
- if(signal(SIGPIPE,lostReader)==SIG_ERR)
- perror("SIGPIPE not supported.");
- pthread_t *pid(new pthread_t[nb_threads]);
- struct pollfd* fds(new struct pollfd[nb_threads]);
- registerMessage(1,testMessage::create);
- registerMessage(2,lastMessage::create);
- int *fdp(new int[nb_threads*2]);
- for(size_t n(0);n<nb_threads;++n){
- if(0>pipe2(&fdp[n*2],0)) {
- perror("pipe thread.");
- exit(-1);
- }
- pthread_create(&pid[n] ,nullptr,&firstThread,&fdp[n*2+1]);
- fds[n].fd=fdp[n*2];
- fds[n].events= POLLIN;
- }
- struct timeval tv,tve;
- gettimeofday(&tv,nullptr);
- int quit(0);
- for(;quit<nb_threads;){
- switch(poll(fds,nb_threads,-1)){
- case 0:
- std::cout << "timeout" << std::endl;
- break;
- case -1:
- perror("Poll");
- exit(-1);
- break;
- default:
- for(size_t n(0);n<nb_threads;++n){
- bool eot(false);
- handleFds(fds[n],ct,c,eot);
- if(eot){
- std::cout << "Last Message from " << fds[n].fd << std::endl;
- ++quit;
- }
- }
- }
- }
- for(size_t n(0);n<nb_threads;++n){
- double *r_thread;
- pthread_join(pid[n],reinterpret_cast<void**>(&r_thread));
- if(r_thread){
- std::cout << "time to write thread " << n <<":\t" << *r_thread << " sec."<< std::endl;
- delete r_thread;
- }
- close(fdp[n*2]);
- close(fdp[n*2+1]);
- }
- delete[] fdp;
- delete[] pid;
- delete[] fds;
- gettimeofday(&tve,nullptr);
- std::cout << "time to read " << c << "/" << ct<< " messages:" << timediff(tv,tve)<< " sec." << std::endl;
- }
|