nitro2005

zeromq_srv_single.cpp

Oct 31st, 2014
374
0
Never
Not a member of Pastebin yet? Sign Up, it unlocks many cool features!
C++ 6.63 KB | None | 0 0
  1. #include "zmq.h"
  2. #include <time.h>
  3. #include <signal.h>
  4. #include <stdio.h>
  5. #include <stdlib.h>
  6. #include <unistd.h>
  7. #include <math.h>
  8. #include <pthread.h>
  9. #include <sys/wait.h>
  10.  
  11. #define SOCKET_STRING                 "tcp://127.0.0.1:1102"
  12. #define PACKET_SIZE                   1024
  13. #define OUTPUT_BUFFER_SIZE            PACKET_SIZE*5
  14. #define MAX_CLIENTS                   10
  15. #define TEST_DURATION                 10
  16.  
  17. struct ThreadStats
  18. {
  19.   unsigned long long packets_received;
  20.   unsigned long long bytes_received;
  21. };
  22.  
  23. bool flag_kill;
  24. void *z_ctx;
  25.  
  26. //---------------------------------------------------------------------------
  27. static void *client_thread(void *)
  28. {
  29.   void *sock = zmq_socket(z_ctx, ZMQ_DEALER);
  30.   if (!sock)
  31.     return NULL;
  32.  
  33.   if (zmq_connect(sock, SOCKET_STRING) < 0)
  34.   {
  35.     zmq_close(sock);
  36.     return NULL;
  37.   }
  38.  
  39.   char output_buffer[OUTPUT_BUFFER_SIZE];
  40.   int output_buffer_len = 0;
  41.  
  42.   // отправляем на сервер несколько пакетов
  43.   while (output_buffer_len+PACKET_SIZE <= sizeof(output_buffer))
  44.   {
  45.     int bytes_sent = zmq_send(sock, output_buffer+output_buffer_len, PACKET_SIZE, ZMQ_DONTWAIT);
  46.     output_buffer_len += bytes_sent;
  47.   }
  48.  
  49.   ThreadStats *stats = new ThreadStats;
  50.   stats->packets_received = 0;
  51.   stats->bytes_received = 0;
  52.  
  53.   // начинаем мониторить события сокета
  54.   zmq_pollitem_t poll_fd;
  55.   poll_fd.socket = sock;
  56.   poll_fd.events = ZMQ_POLLIN;
  57.   poll_fd.revents = 0;
  58.  
  59.   while(!flag_kill)
  60.   {
  61.     int res = zmq_poll(&poll_fd, 1, 500);
  62.     if (res < 0)
  63.       break;
  64.     if (res == 0)
  65.       continue;
  66.  
  67.     // есть новый пакет
  68.     if (poll_fd.revents & ZMQ_POLLIN)
  69.     {
  70.       poll_fd.revents &= ~ZMQ_POLLIN;
  71.  
  72.       // читаем пакет с данными
  73.       zmq_msg_t rcv_msg;
  74.       if (zmq_msg_init(&rcv_msg) < 0)
  75.       {
  76.         flag_kill = true;
  77.         break;
  78.       }
  79.  
  80.       int rc = zmq_msg_recv(&rcv_msg, sock, ZMQ_DONTWAIT);
  81.       if (rc < 0)
  82.       {
  83.         flag_kill = true;
  84.         break;
  85.       }
  86.  
  87.       int msg_size = zmq_msg_size(&rcv_msg);
  88.       char *msg_buffer = (char *)zmq_msg_data(&rcv_msg);
  89.       stats->bytes_received += msg_size;
  90.       stats->packets_received++;
  91.  
  92.       // отправляем пакет с данными
  93.       if (zmq_send(sock, msg_buffer, msg_size, ZMQ_DONTWAIT) < 0)
  94.       {
  95.         printf("zmq_send() error: %s\r\n", zmq_strerror(errno));
  96.         break;
  97.       }
  98.  
  99.       zmq_msg_close(&rcv_msg);
  100.     }
  101.   }
  102.  
  103.   zmq_close(sock);
  104.   return stats;
  105. }
  106.  
  107. //---------------------------------------------------------------------------
  108. void termination_handler(int)
  109. {
  110.   flag_kill = true;
  111. }
  112. //---------------------------------------------------------------------------
  113. int main(int argc, char *argv[])
  114. {
  115.   signal(SIGTERM, termination_handler);
  116.   signal(SIGSTOP, termination_handler);
  117.   signal(SIGINT,  termination_handler);
  118.   signal(SIGQUIT, termination_handler);
  119.  
  120.   z_ctx = zmq_ctx_new();
  121.   if (!z_ctx)
  122.     return 1;
  123.  
  124. //  zmq_ctx_set(z_ctx, ZMQ_IO_THREADS, 2);
  125.  
  126.   void *sock_srv = zmq_socket(z_ctx, ZMQ_ROUTER);
  127.   if (!sock_srv)
  128.   {
  129.     zmq_ctx_destroy(z_ctx);
  130.     return 2;
  131.   }
  132.  
  133.   if (zmq_bind(sock_srv, SOCKET_STRING) < 0)
  134.   {
  135.     zmq_close(sock_srv);
  136.     zmq_ctx_destroy(z_ctx);
  137.     return 3;
  138.   }
  139.  
  140.   // запускаем потоки клиентов
  141.   pthread_t thread_ids[MAX_CLIENTS];
  142.   for (int i=0; i < MAX_CLIENTS; i++)
  143.   {
  144.     pthread_create(&thread_ids[i], NULL, &client_thread, NULL);
  145.   }
  146.  
  147.   // начинаем мониторить события сокета
  148.   zmq_pollitem_t poll_fd;
  149.   poll_fd.socket = sock_srv;
  150.   poll_fd.events = ZMQ_POLLIN;
  151.   poll_fd.revents = 0;
  152.  
  153.   struct timespec ts_start;
  154.   struct timespec ts_current;
  155.   clock_gettime(CLOCK_MONOTONIC, &ts_start);
  156.   double start_time = ts_start.tv_sec + (double)ts_start.tv_nsec/1000000000;
  157.   double cur_time = ts_start.tv_sec + (double)ts_start.tv_nsec/1000000000;
  158.   zmq_msg_t rcv_msg;
  159.   while(!flag_kill)
  160.   {
  161.     clock_gettime(CLOCK_MONOTONIC, &ts_current);
  162.     cur_time = ts_current.tv_sec + (double)ts_current.tv_nsec/1000000000;
  163.     if (cur_time-start_time > TEST_DURATION)
  164.     {
  165.       flag_kill = true;
  166.       break;
  167.     }
  168.     int res = zmq_poll(&poll_fd, 1, 500);
  169.     if (res < 0)
  170.       break;
  171.     if (res == 0)
  172.       continue;
  173.  
  174.     // есть новый пакет
  175.     if (poll_fd.revents & ZMQ_POLLIN)
  176.     {
  177.       poll_fd.revents &= ~ZMQ_POLLIN;
  178.  
  179.       // читаем фрейм с идентификатором клиента
  180.       char client_identity[255];
  181.       int client_identity_len;
  182.       client_identity_len = zmq_recv(sock_srv, client_identity, sizeof(client_identity), ZMQ_DONTWAIT);
  183.       if (client_identity_len <= 0)
  184.         break;
  185.  
  186.       // читаем фрейм с пакетом данных
  187.       if (zmq_msg_init(&rcv_msg) < 0)
  188.       {
  189.         flag_kill = true;
  190.         break;
  191.       }
  192.       int rc = zmq_msg_recv(&rcv_msg, sock_srv, ZMQ_DONTWAIT);
  193.       if (rc < 0)
  194.       {
  195.         flag_kill = true;
  196.         break;
  197.       }
  198.       int msg_size = zmq_msg_size(&rcv_msg);
  199.  
  200.       // отправляем фрейм с идентификатором клиента
  201.       if (zmq_send(sock_srv, client_identity, client_identity_len, ZMQ_DONTWAIT | ZMQ_SNDMORE) < 0)
  202.         break;
  203.  
  204.       // отправляем фрейм с пакетом данных
  205.       if (zmq_send(sock_srv, zmq_msg_data(&rcv_msg), msg_size, ZMQ_DONTWAIT) < 0)
  206.         break;
  207.  
  208.       zmq_msg_close(&rcv_msg);
  209.     }
  210.   }
  211.   double elapsed = cur_time-start_time;
  212.  
  213.   // ждем завершения потоков
  214.   usleep(10000);
  215.  
  216.   unsigned long long total_packets_received = 0;
  217.   unsigned long long total_bytes_received = 0;
  218.   // получаем статистику по потокам и считаем итог
  219.   printf("thread id:\tbytes rcv\tpackets rcv\r\n");
  220.   for (int i=0; i < MAX_CLIENTS; i++)
  221.   {
  222.     void *res = NULL;
  223.     pthread_join(thread_ids[i], &res);
  224.     ThreadStats *stat = (ThreadStats *)res;
  225.     if (stat)
  226.     {
  227.       printf("thread %02d:\t%lld\t%lld\r\n", i+1, stat->bytes_received, stat->packets_received);
  228.       total_bytes_received += stat->bytes_received;
  229.       total_packets_received += stat->packets_received;
  230.       delete stat;
  231.     }
  232.   }
  233.   printf("    TOTAL:\t%lld\t%lld\r\n", total_bytes_received, total_packets_received);
  234.   printf("\r\n");
  235.   printf("Elapsed time: %.3lf\r\n", elapsed);
  236.   printf("Avg speed: %d bytes/s, %d packets/s\r\n", (unsigned int)(total_bytes_received/elapsed), (unsigned int)(total_packets_received/elapsed));
  237.  
  238.   zmq_close(sock_srv);
  239.   zmq_ctx_destroy(z_ctx);
  240.   return 0;
  241. }
Advertisement
Add Comment
Please, Sign In to add comment