nitro2005

socket_srv_single.cpp

Oct 31st, 2014
325
0
Never
Not a member of Pastebin yet? Sign Up, it unlocks many cool features!
C++ 14.43 KB | None | 0 0
  1. #include <unistd.h>
  2. #include <sys/types.h>
  3. #include <sys/wait.h>
  4. #include <signal.h>
  5. #include <stdio.h>
  6. #include <math.h>
  7. #include <string.h>
  8. #include <stdlib.h>
  9. #include <time.h>
  10. #include <malloc.h>
  11. #include <fcntl.h>
  12. #include <syslog.h>
  13. #include <sys/socket.h>
  14. #include <netinet/in.h>
  15. #include <arpa/inet.h>
  16. #include <netdb.h>
  17. #include <pthread.h>
  18. #include <sys/poll.h>
  19.  
  20. #define SOCKET_TCP_PORT               1101
  21. #define PACKET_SIZE                   1024
  22. #define INPUT_BUFFER_SIZE             (PACKET_SIZE+sizeof(int))*5
  23. #define OUTPUT_BUFFER_SIZE            (PACKET_SIZE+sizeof(int))*5
  24. #define MAX_CLIENTS                   10
  25. #define TEST_DURATION                 10
  26.  
  27. struct SocketClient
  28. {
  29.   char input_buffer[INPUT_BUFFER_SIZE];
  30.   char output_buffer[OUTPUT_BUFFER_SIZE];
  31.   int input_buffer_len;
  32.   int output_buffer_len;
  33.   int fd;
  34.   int poll_fd_index;
  35. };
  36. struct ThreadStats
  37. {
  38.   unsigned long long packets_received;
  39.   unsigned long long bytes_received;
  40. };
  41.  
  42. struct SocketClient clients[MAX_CLIENTS];
  43. bool flag_kill;
  44.  
  45. //---------------------------------------------------------------------------
  46. bool setnonblocking(int fd)
  47. {
  48.   int opts = fcntl(fd, F_GETFL);
  49.   if (opts < 0)
  50.     return false;
  51.   opts = (opts | O_NONBLOCK);
  52.   if (fcntl(fd, F_SETFL, opts) < 0)
  53.     return false;
  54.   return true;
  55. }
  56.  
  57. //---------------------------------------------------------------------------
  58. static void *client_thread(void *)
  59. {
  60.   // устанавливаем соединение с сервером
  61.   struct hostent *host_info = gethostbyname("127.0.0.1");
  62.   if (!host_info)
  63.     return NULL;
  64.   int fd = socket(AF_INET, SOCK_STREAM, 0);
  65.   if (fd < 0)
  66.     return NULL;
  67.   struct sockaddr_in addr;
  68.   memset((char *) &addr, 0, sizeof(addr));
  69.   addr.sin_family = host_info->h_addrtype;
  70.   addr.sin_port = htons(SOCKET_TCP_PORT);
  71.   addr.sin_addr = *(struct in_addr*)host_info->h_addr;
  72.   if (connect(fd, (struct sockaddr *)&addr, sizeof(addr)) < 0)
  73.   {
  74.     close(fd);
  75.     return NULL;
  76.   }
  77.   if (!setnonblocking(fd))
  78.   {
  79.     close(fd);
  80.     return NULL;
  81.   }
  82.  
  83.   // начинаем мониторить события сокета
  84.   struct pollfd poll_fd;
  85.   memset(&poll_fd, 0, sizeof(poll_fd));
  86.   poll_fd.fd = fd;
  87.   poll_fd.events = POLLIN | POLLOUT | POLLERR | POLLHUP;
  88.   poll_fd.revents = 0;
  89.  
  90.   char input_buffer[INPUT_BUFFER_SIZE];
  91.   char output_buffer[OUTPUT_BUFFER_SIZE];
  92.   int input_buffer_len = 0;
  93.   int output_buffer_len = 0;
  94.  
  95.   // кладем в буфер на передачу несколько пакетов
  96.   while (output_buffer_len+sizeof(int)+PACKET_SIZE <= sizeof(output_buffer))
  97.   {
  98.     *((int *)(output_buffer+output_buffer_len)) = PACKET_SIZE;
  99.     output_buffer_len += sizeof(int)+PACKET_SIZE;
  100.   }
  101.  
  102.   ThreadStats *stats = new ThreadStats;
  103.   stats->bytes_received = 0;
  104.   stats->packets_received = 0;
  105.  
  106.   while(!flag_kill)
  107.   {
  108.     // если есть что передавать - мониторим разрешение на передачу
  109.     if (output_buffer_len > 0)
  110.       poll_fd.events |= POLLOUT;
  111.     else
  112.       poll_fd.events &= ~POLLOUT;
  113.  
  114.     int res = poll(&poll_fd, 1, 500);
  115.     if (res < 0)
  116.       break;
  117.     if (res == 0)
  118.       continue;
  119.  
  120.     // ошибочка
  121.     if ((poll_fd.revents & POLLERR) || (poll_fd.revents & POLLHUP))
  122.       break;
  123.  
  124.     // есть данные для чтения из сокета
  125.     if (poll_fd.revents & POLLIN)
  126.     {
  127.       int rcv_len = recv(fd, input_buffer+input_buffer_len, sizeof(input_buffer)-input_buffer_len, 0);
  128.       if (rcv_len <= 0)
  129.         break;
  130.       stats->bytes_received += rcv_len;
  131.       input_buffer_len += rcv_len;
  132.  
  133.       // парсинг буфера приема
  134.       while(input_buffer_len > sizeof(int))
  135.       {
  136.         int packet_size = *((int *)input_buffer);
  137.         if (input_buffer_len < sizeof(int)+packet_size)
  138.           break;
  139.         if (output_buffer_len+sizeof(int)+packet_size > sizeof(output_buffer))
  140.           break;
  141.         stats->packets_received++;
  142.         // если получили корректный пакет и можем его впихнуть в буфер на передачу - отправляем его туда
  143.         memcpy(output_buffer+output_buffer_len, input_buffer, sizeof(int)+packet_size);
  144.         output_buffer_len += sizeof(int)+packet_size;
  145.         // а из буфера приема удаляем
  146.         memmove(input_buffer, input_buffer+sizeof(int)+packet_size, input_buffer_len-sizeof(int)-packet_size);
  147.         input_buffer_len -= sizeof(int)+packet_size;
  148.       }
  149.  
  150.       poll_fd.revents &= ~POLLIN;
  151.     }
  152.  
  153.     // можно записать данные в сокет
  154.     if (poll_fd.revents & POLLOUT)
  155.     {
  156.       int snd_len = send(fd, output_buffer, output_buffer_len, 0);
  157.       if (snd_len <= 0)
  158.         break;
  159.       // удаляем отправленные данные из буфера
  160.       memmove(output_buffer, output_buffer+snd_len, output_buffer_len-snd_len);
  161.       output_buffer_len -= snd_len;
  162.       poll_fd.revents &= ~POLLOUT;
  163.     }
  164.   }
  165.   close(fd);
  166.  
  167.   return stats;
  168. }
  169.  
  170. //---------------------------------------------------------------------------
  171. void termination_handler(int)
  172. {
  173.   flag_kill = true;
  174. }
  175. //---------------------------------------------------------------------------
  176. void client_parse_input_buffer(int client_index)
  177. {
  178.   while(clients[client_index].input_buffer_len > sizeof(int))
  179.   {
  180.     int packet_size = *((int *)clients[client_index].input_buffer);
  181.     if (clients[client_index].input_buffer_len < sizeof(int)+packet_size)
  182.       break;
  183.     if (clients[client_index].output_buffer_len+sizeof(int)+packet_size > sizeof(clients[client_index].output_buffer))
  184.       break;
  185.     // если получили корректный пакет и можем его впихнуть в буфер на передачу - отправляем его туда
  186.     memcpy(clients[client_index].output_buffer+clients[client_index].output_buffer_len, clients[client_index].input_buffer, sizeof(int)+packet_size);
  187.     clients[client_index].output_buffer_len += sizeof(int)+packet_size;
  188.     // а из буфера приема удаляем
  189.     memmove(clients[client_index].input_buffer, clients[client_index].input_buffer+sizeof(int)+packet_size, clients[client_index].input_buffer_len-sizeof(int)-packet_size);
  190.     clients[client_index].input_buffer_len -= sizeof(int)+packet_size;
  191.   }
  192. }
  193. //---------------------------------------------------------------------------
  194. void client_disconnected(int client_index, struct pollfd *poll_fds, int &n_fds)
  195. {
  196.   close(clients[client_index].fd);
  197.   clients[client_index].fd = 0;
  198.   if (clients[client_index].poll_fd_index >= 0 && clients[client_index].poll_fd_index < n_fds-1)
  199.   {
  200.     memmove(poll_fds+clients[client_index].poll_fd_index,
  201.             poll_fds+clients[client_index].poll_fd_index+1,
  202.             (n_fds-clients[client_index].poll_fd_index-1)*sizeof(struct pollfd));
  203.   }
  204.   n_fds--;
  205. }
  206. //---------------------------------------------------------------------------
  207. int main(int argc, char *argv[])
  208. {
  209.   signal(SIGTERM, termination_handler);
  210.   signal(SIGSTOP, termination_handler);
  211.   signal(SIGINT,  termination_handler);
  212.   signal(SIGQUIT, termination_handler);
  213.  
  214.   int sock_srv = socket(AF_INET, SOCK_STREAM, 0);
  215.   if (sock_srv < 0)
  216.     return 1;
  217.  
  218.   // устанавливаем опцию для повторного использования адреса (порта) без ожидания таймаута
  219.   int reuse_addr = 1;
  220.   if (setsockopt(sock_srv, SOL_SOCKET, SO_REUSEADDR, &reuse_addr, sizeof(reuse_addr)) < 0)
  221.     return 2;
  222.  
  223.   // устанавливаем неблокирующий режим
  224.   if (!setnonblocking(sock_srv))
  225.     return 3;
  226.  
  227.   struct sockaddr_in addr;
  228.   memset((char *) &addr, 0, sizeof(addr));
  229.   addr.sin_family = AF_INET;
  230.   addr.sin_port = htons(SOCKET_TCP_PORT);
  231.   addr.sin_addr.s_addr = htonl(INADDR_ANY);
  232.  
  233.   // привязываем сокет к адресу (порту)
  234.   if (bind(sock_srv, (struct sockaddr *)&addr, sizeof(addr)) < 0)
  235.     return 4;
  236.  
  237.   // начинаем слушать порт
  238.   if (listen(sock_srv, 100) < 0)
  239.     return 5;
  240.  
  241.   // запускаем потоки клиентов
  242.   pthread_t thread_ids[MAX_CLIENTS];
  243.   for (int i=0; i < MAX_CLIENTS; i++)
  244.   {
  245.     pthread_create(&thread_ids[i], NULL, &client_thread, NULL);
  246.   }
  247.  
  248.   // чистим массив с данными клиентов
  249.   for (int i=0; i < MAX_CLIENTS; i++)
  250.   {
  251.     clients[i].fd = -1;
  252.     clients[i].input_buffer_len = 0;
  253.     clients[i].output_buffer_len = 0;
  254.     clients[i].poll_fd_index = -1;
  255.   }
  256.  
  257.   // начинаем мониторить события сокета(ов)
  258.   // по мере подключения клиентов n_fds будет возрастать, poll_fds заполняться
  259.   // а отключение первого же клиента будет означать останов программы,
  260.   // поэтому мы не будем заботиться о сдвиге poll_fds при отключении клиента
  261.   struct pollfd poll_fds[MAX_CLIENTS+1];
  262.   memset(&poll_fds, 0, sizeof(pollfd)*(MAX_CLIENTS+1));
  263.   poll_fds[0].fd = sock_srv;
  264.   poll_fds[0].events = POLLIN;
  265.   poll_fds[0].revents = 0;
  266.   int n_fds = 1;
  267.  
  268.   struct timespec ts_start;
  269.   struct timespec ts_current;
  270.   clock_gettime(CLOCK_MONOTONIC, &ts_start);
  271.   double start_time = ts_start.tv_sec + (double)ts_start.tv_nsec/1000000000;
  272.   double cur_time = ts_start.tv_sec + (double)ts_start.tv_nsec/1000000000;
  273.   while(!flag_kill)
  274.   {
  275.     clock_gettime(CLOCK_MONOTONIC, &ts_current);
  276.     cur_time = ts_current.tv_sec + (double)ts_current.tv_nsec/1000000000;
  277.     if (cur_time-start_time > TEST_DURATION)
  278.     {
  279.       flag_kill = true;
  280.       break;
  281.     }
  282.     // если есть что передавать - мониторим разрешение на передачу
  283.     for (int i=0; i < MAX_CLIENTS; i++)
  284.     {
  285.       if (clients[i].fd >= 0)
  286.       {
  287.         if (clients[i].output_buffer_len > 0)
  288.           poll_fds[clients[i].poll_fd_index].events |= POLLOUT;
  289.         else
  290.           poll_fds[clients[i].poll_fd_index].events &= ~POLLOUT;
  291.       }
  292.     }
  293.  
  294.     int res = poll(poll_fds, n_fds, 500);
  295.     if (res < 0)
  296.       break;
  297.     if (res == 0)
  298.       continue;
  299.  
  300.     // есть новое подключение клиента
  301.     if (poll_fds[0].revents & POLLIN)
  302.     {
  303.       poll_fds[0].revents &= ~POLLIN;
  304.       socklen_t len = sizeof(struct sockaddr_in);
  305.       struct sockaddr_in addr;
  306.       memset(&addr, 0, len);
  307.       int sock = accept(sock_srv, (struct sockaddr *)&addr, &len);
  308.       if (sock < 0)
  309.         break;
  310.       if (!setnonblocking(sock))
  311.         break;
  312.       // если слишком много клиентов - закрываем новое соединение
  313.       if (n_fds > MAX_CLIENTS)
  314.       {
  315.         close(sock);
  316.       }
  317.       else
  318.       {
  319.         int client_index = -1;
  320.         for (int i=0; i < MAX_CLIENTS; i++)
  321.         {
  322.           if (clients[i].fd < 0)
  323.           {
  324.             client_index = i;
  325.             break;
  326.           }
  327.         }
  328.         if (client_index >= 0)
  329.         {
  330.           clients[client_index].fd = sock;
  331.           clients[client_index].input_buffer_len = 0;
  332.           clients[client_index].output_buffer_len = 0;
  333.           clients[client_index].poll_fd_index = n_fds;
  334.           poll_fds[n_fds].fd = sock;
  335.           poll_fds[n_fds].events = POLLIN | POLLERR | POLLHUP;
  336.           poll_fds[n_fds].revents = 0;
  337.           n_fds++;
  338.         }
  339.         else
  340.           close(sock);
  341.       }
  342.     }
  343.  
  344.     for (int i=0; i < MAX_CLIENTS; i++)
  345.     {
  346.       if (clients[i].fd >= 0)
  347.       {
  348.         if ((poll_fds[clients[i].poll_fd_index].revents & POLLERR) || (poll_fds[clients[i].poll_fd_index].revents & POLLHUP))
  349.         {
  350.           client_disconnected(i, poll_fds, n_fds);
  351.           continue;
  352.         }
  353.  
  354.         // если есть что читать от клиента
  355.         if (poll_fds[clients[i].poll_fd_index].revents & POLLIN)
  356.         {
  357.           poll_fds[clients[i].poll_fd_index].revents &= ~POLLIN;
  358.           int rcv_len = recv(clients[i].fd, clients[i].input_buffer+clients[i].input_buffer_len, sizeof(clients[i].input_buffer)-clients[i].input_buffer_len, 0);
  359.           if (rcv_len <= 0)
  360.           {
  361.             client_disconnected(i, poll_fds, n_fds);
  362.             continue;
  363.           }
  364.           clients[i].input_buffer_len += rcv_len;
  365.  
  366.           // парсинг буфера приема
  367.           client_parse_input_buffer(i);
  368.         }
  369.  
  370.         // если можно записать данные в сокет
  371.         if (poll_fds[clients[i].poll_fd_index].revents & POLLOUT)
  372.         {
  373.           poll_fds[clients[i].poll_fd_index].revents &= ~POLLOUT;
  374.           int snd_len = send(clients[i].fd, clients[i].output_buffer, clients[i].output_buffer_len, 0);
  375.           if (snd_len <= 0)
  376.           {
  377.             client_disconnected(i, poll_fds, n_fds);
  378.             continue;
  379.           }
  380.           // удаляем отправленные данные из буфера
  381.           memmove(clients[i].output_buffer, clients[i].output_buffer+snd_len, clients[i].output_buffer_len-snd_len);
  382.           clients[i].output_buffer_len -= snd_len;
  383.         }
  384.       }
  385.     }
  386.   }
  387.   double elapsed = cur_time-start_time;
  388.  
  389.   // отключаем клиентов
  390.   for (int i=0; i < MAX_CLIENTS; i++)
  391.   {
  392.     if (clients[i].fd >= 0)
  393.       close(clients[i].fd);
  394.   }
  395.  
  396.   // ждем завершения потоков
  397.   usleep(10000);
  398.  
  399.   unsigned long long total_packets_received = 0;
  400.   unsigned long long total_bytes_received = 0;
  401.   // запускаем потоки клиентов
  402.   printf("thread id:\tbytes rcv\tpackets rcv\r\n");
  403.   for (int i=0; i < MAX_CLIENTS; i++)
  404.   {
  405.     void *res = NULL;
  406.     pthread_join(thread_ids[i], &res);
  407.     ThreadStats *stat = (ThreadStats *)res;
  408.     if (stat)
  409.     {
  410.       printf("thread %02d:\t%lld\t%lld\r\n", i+1, stat->bytes_received, stat->packets_received);
  411.       total_bytes_received += stat->bytes_received;
  412.       total_packets_received += stat->packets_received;
  413.       delete stat;
  414.     }
  415.   }
  416.   printf("    TOTAL:\t%lld\t%lld\r\n", total_bytes_received, total_packets_received);
  417.   printf("\r\n");
  418.   printf("Elapsed time: %.3lf\r\n", elapsed);
  419.   printf("Avg speed: %d bytes/s, %d packets/s\r\n", (unsigned int)(total_bytes_received/elapsed), (unsigned int)(total_packets_received/elapsed));
  420.  
  421.   close(sock_srv);
  422.   return 0;
  423. }
Advertisement
Add Comment
Please, Sign In to add comment