Guest User

Untitled

a guest
Oct 29th, 2022
188
0
Never
Not a member of Pastebin yet? Sign Up, it unlocks many cool features!
C 14.91 KB | None | 0 0
  1. /**
  2.  * \file src/sstrx_qrecord.c
  3.  * \brief Main module for quick timeseries recorder
  4.  * \author *
  5.  *
  6.  * Implements a quick timeseries data recorder client for sstrx
  7.  *
  8.  */
  9.  
  10. #include "config.h"
  11.  
  12. #include <stdio.h>
  13. #include <math.h>
  14. #include <syslog.h>
  15. #include <stdlib.h>
  16. #include <string.h>
  17. #include <unistd.h>
  18. #include <getopt.h>
  19. #include <pthread.h>
  20. #include <signal.h>
  21. #include <sys/socket.h>
  22. #include <sys/types.h>
  23. #include <sys/stat.h>
  24. #include <fcntl.h>
  25. #include <sys/types.h>
  26. #include <sys/time.h>
  27. #include <errno.h>
  28. #include <time.h>
  29. #include <poll.h>
  30. #include <limits.h>
  31.  
  32. #include "net_utils.h"
  33. #include "client_lib.h"
  34. #include "data_format.h"
  35. #include "radar_info.h"
  36. #include "queue.h"
  37.  
  38. /* * * * Module-local constants * * * */
  39. enum PROGRAM_ARGS {
  40.     ARG_HELP = 256,
  41.     ARG_SERVERNAME,
  42.     ARG_TIMEOUT,
  43.     ARG_SAMPLES,
  44.     ARG_QUEUESIZE,
  45.     ARG_NUM_PRTS,
  46.     ARG_OUTPUT_DIR,
  47.     ARG_FLUSHAFTER,
  48. };
  49. #define MAX_HOSTNAME 128
  50. #define MAX_FILENAME PATH_MAX
  51. #define MAX_PATHNAME PATH_MAX
  52.  
  53. #define DEFAULT_QUEUE_SIZE 2000
  54. #define DEFAULT_SAMPLES 5000
  55. #define DEFAULT_FLUSHAFTER_MB 32
  56.  
  57. /* * * * Module-local structure definitions * * * */
  58. struct qrecord_data {
  59.     uint32_t num_samples;
  60.     char filename[MAX_FILENAME];
  61.     int filehandle;
  62.     struct queue queue;
  63.     int buffer_size;
  64.     int queue_size;
  65.     int num_prts;
  66.     int show_timestamps;
  67.     uint32_t automatic_mode;
  68.     uint64_t bytes_written;
  69.     char output_path[MAX_PATHNAME];
  70.     int queue_full;
  71.     int pulses_dropped;
  72.     uint64_t flush_after_bytes;
  73.     uint64_t bytes_since_flush;
  74. };
  75.  
  76. /* * * * Module-local function declarations * * * */
  77. static int do_qrecord(struct sstrx_acqd *sa, struct qrecord_data *qr);
  78. static void setup_signals(void);
  79. static int process_header(union sstrx_acqd_cmd *hdr, struct qrecord_data *qr);
  80. static int process_data(struct sstrx_acqd_data_header *hdr, struct qrecord_data *qr);
  81. void *disk_write_thread(void *argument);
  82.  
  83. /* * * * Module-local variables * * * */
  84. static char hostname[MAX_HOSTNAME] = "localhost";
  85. static uint32_t port = SSTRX_ACQD_DEFAULT_PORT;
  86. volatile static int g_terminate_flag = 0;
  87. static struct qrecord_data qrecord_data = {
  88.     .num_samples = DEFAULT_SAMPLES,
  89.     .num_prts = 0,
  90.     .queue_size = DEFAULT_QUEUE_SIZE,
  91.     .flush_after_bytes = 1048576 * DEFAULT_FLUSHAFTER_MB,
  92. };
  93.  
  94. /** Usage string, printed when user needs help w/ command line */
  95. static char *usage =
  96. "SSTRX quick timeseries recorder\n"
  97. "sstrx-acqd-qrecord [options] [filename]\n"
  98. " [filename]: File to record to.\n"
  99. "             If no file is specified, enter automatic mode where the recorder\n"
  100. "             waits for event notifications to start and stop recording, files\n"
  101. "             are created with the names TS_yyyymmdd_hhmmss.dat"
  102. "--help: Show usage information\n"
  103. "--server=<hostname:port>: Set the server name and port. Defaults to\n"
  104. "                          env. var. ACQD_SVR if available, else localhost\n"
  105. "--timeout=<timeout>: Set the timeout in seconds (fractional OK)\n"
  106. "--samples=<samples>: Acquire <samples> samples per PRT, defaults to " xstr(DEFAULT_SAMPLES) "\n"
  107. "--queuesize=<size>: Set queue size to <size>, defaults to " xstr(DEFAULT_QUEUE_SIZE) "\n"
  108. "--prts=<prts>: Stop recording after <prts> PRTs.\n"
  109. "--showtimestamps: Display timestamp information\n"
  110. "--outdir=<path>: Write to specified directory (default: current dir)\n"
  111. "--flushafter=<mbytes>: Flush cache after specified megabytes, default " xstr(DEFAULT_FLUSHAFTER_MB) "M \n"
  112. "\n";
  113.  
  114. /* * * * Function definitions * * * */
  115.  
  116. int main(int argc, char *argv[])
  117. {
  118.     int opt, opt_idx, opt_errors = 0;
  119.     int rv = -1;
  120.     struct option opt_lst[] = {
  121.         {"help", no_argument, 0, ARG_HELP},
  122.         {"server", required_argument, 0, ARG_SERVERNAME},
  123.         {"timeout", required_argument, 0, ARG_TIMEOUT},
  124.         {"samples", required_argument, 0, ARG_SAMPLES},
  125.         {"queuesize", required_argument, 0, ARG_QUEUESIZE},
  126.         {"prts", required_argument, 0, ARG_NUM_PRTS},
  127.         {"showtimestamps", no_argument, &qrecord_data.show_timestamps, 1},
  128.         {"outdir", required_argument, 0, ARG_OUTPUT_DIR},
  129.         {"flushafter", required_argument, 0, ARG_FLUSHAFTER},
  130.         {0, 0, 0, 0} // Sentinel
  131.     };
  132.  
  133.     struct sstrx_acqd sstrx_acqd;
  134.     struct sstrx_acqd *sa = &sstrx_acqd;
  135.  
  136.     if (NULL != getenv("ACQD_SVR")) {
  137.         strncpy(hostname, getenv("ACQD_SVR"), MAX_HOSTNAME-1);
  138.         strtok(hostname, ":");
  139.         char *port_str = strtok(NULL, ":");
  140.         port = (port_str != 0) ? strtoul(port_str, NULL, 10) : SSTRX_ACQD_DEFAULT_PORT;
  141.     }
  142.  
  143.     sa->timeout = 1000;
  144.     sa->max_skip = 100;
  145.     sa->input_cmd_buffer_size = sstrx_acqd_buffer_size(SSTRX_ACQD_BUFSZ_ALL);
  146.     //printf("Allocating buffer of size %d bytes\n", sa->input_cmd_buffer_size);
  147.     sa->input_cmd_buffer = malloc(sa->input_cmd_buffer_size);
  148.  
  149.     struct qrecord_data *qr = &qrecord_data;
  150.  
  151.     while (-1 != (opt = getopt_long(argc, argv, "", opt_lst, &opt_idx))) {
  152.         switch (opt) {
  153.         case ARG_HELP:
  154.             puts(usage);
  155.             return 0;
  156.             break;
  157.         case ARG_SERVERNAME:
  158.             {
  159.             strncpy(hostname, strtok(optarg, ":"), MAX_HOSTNAME-1);
  160.             char *port_str = strtok(NULL, ":");
  161.             port = (port_str != 0) ? strtoul(port_str, NULL, 10) : SSTRX_ACQD_DEFAULT_PORT;
  162.             }
  163.             break;
  164.         case ARG_TIMEOUT:
  165.             sa->timeout = round(strtod(optarg, NULL) * 1e3);
  166.             break;
  167.         case ARG_SAMPLES:
  168.             qr->num_samples = strtoul(optarg, NULL, 10);
  169.             break;
  170.         case ARG_QUEUESIZE:
  171.             qr->queue_size = strtoul(optarg, NULL, 10);
  172.             break;
  173.         case ARG_NUM_PRTS:
  174.             qr->num_prts = strtoul(optarg, NULL, 10);
  175.             break;
  176.         case ARG_OUTPUT_DIR:
  177.             stpncpy(qr->output_path, optarg, MAX_PATHNAME-1);
  178.             strncat(qr->output_path, "/", MAX_PATHNAME);
  179.             break;
  180.         case ARG_FLUSHAFTER:
  181.             qr->flush_after_bytes = strtoul(optarg, NULL, 10) * 1048576;
  182.             break;
  183.         case 0:
  184.             break;
  185.         default:
  186.             opt_errors++;
  187.             break;
  188.         }
  189.     }
  190.  
  191.     qr->buffer_size = sa->input_cmd_buffer_size;
  192.     if (0 != queue_init(&qr->queue, qr->queue_size, sa->input_cmd_buffer_size)) {
  193.         fprintf(stderr, "Queue init failed\n");
  194.         goto terminate_with_error;
  195.     }
  196.  
  197.     if (opt_errors) {
  198.         puts(usage);
  199.         return -1;
  200.     }
  201.  
  202.     if (optind < argc) {
  203.         qr->automatic_mode = 0;
  204.         strncpy(qr->filename, argv[optind], MAX_FILENAME);
  205.         int flags = O_WRONLY | O_CREAT | O_LARGEFILE | O_TRUNC;
  206.         mode_t mode = S_IWUSR | S_IRUSR | S_IRGRP | S_IROTH;
  207.  
  208.         qr->filehandle = open(qr->filename, flags, mode);
  209.         printf("open returns: %d\n", qr->filehandle);
  210.         if (qr->filehandle < 0) {
  211.             perror("open");
  212.         }
  213.     }
  214.     else {
  215.         qr->automatic_mode = 1;
  216.         strcpy(qr->filename, "<none>");
  217.         qr->filehandle = -1;
  218.     }
  219.  
  220.     printf("Connecting to %s:%d\n", hostname, port);
  221.     sa->fd = open_socket(hostname, port);
  222.     if (-1 == sa->fd) {
  223.         fprintf(stderr, "Could not connect to %s:%d\n", hostname, port);
  224.         return -1;
  225.     }
  226.  
  227.     setup_signals();
  228.     rv = sstrx_acqd_set_data_format(sa,
  229.             SSTRX_ACQD_FLAGS_FMT_24B_2CH,
  230.             qr->num_samples, 0);
  231.     if(0 > rv) goto terminate_with_error;
  232.     rv = do_qrecord(sa, qr);
  233.     if(0 > rv) goto terminate_with_error;
  234.  
  235.     if (sa->fd != -1) {
  236.         shutdown(sa->fd, SHUT_RDWR);
  237.         close(sa->fd);
  238.     }
  239.     return 0;
  240.  
  241. terminate_with_error:
  242.     if (sa->fd != -1) {
  243.         close(sa->fd);
  244.     }
  245.     if (0 == rv) {
  246.         printf("No connection to server %s:%d\n", hostname, port);
  247.     }
  248.     else {
  249.         printf("Error: %s\n", enum_to_string(sstrx_acqd_client_errors, rv));
  250.     }
  251.     return rv;
  252.  
  253. }
  254.  
  255. static int do_qrecord(struct sstrx_acqd *sa, struct qrecord_data *qr)
  256. {
  257.     struct sstrx_acqd_data_header *hdr = (struct sstrx_acqd_data_header *)sa->input_cmd_buffer;
  258.     int prts_acquired = 0;
  259.  
  260.     pthread_t dwt_id;
  261.     pthread_create(&dwt_id, NULL, disk_write_thread, qr);
  262.     int bytesread, byteswanted;
  263.     while(!g_terminate_flag) {
  264.         struct pollfd pfd = {sa->fd, POLLIN, 0};
  265.         int rv = poll(&pfd, 1, 1000);
  266.         if (1 != rv) {
  267.             if (g_terminate_flag) break;
  268.             continue;
  269.         }
  270.         byteswanted = sizeof(struct sstrx_acqd_cmd_generic);
  271.         bytesread = safe_read(sa->fd, hdr, byteswanted);
  272.         if (bytesread != byteswanted) {
  273.             fprintf(stderr, "error reading ID/len: %d\n", bytesread);
  274.             goto qrecord_done;
  275.         }
  276.         void *write_location = ((struct sstrx_acqd_cmd_generic *)hdr) + 1;
  277.         byteswanted = hdr->len - sizeof(struct sstrx_acqd_cmd_generic);
  278.         uint32_t temp;
  279.         bytesread = safe_read_max(sa->fd, write_location, byteswanted,
  280.                 sa->input_cmd_buffer_size, &temp, sizeof(temp));
  281.         if (bytesread < 0) {
  282.             perror("read");
  283.             goto qrecord_done;
  284.         }
  285.  
  286.         if (0 != queue_try_add(&qr->queue, sa->input_cmd_buffer, hdr->len)) {
  287.             qr->queue_full = 1;
  288.         }
  289.         else {
  290.             qr->queue_full = 0;
  291.         }
  292.         prts_acquired++;
  293.         if (qr-> num_prts && (prts_acquired >= qr->num_prts)) {
  294.             g_terminate_flag = 1;
  295.             break;
  296.         }
  297.     }
  298.     close(sa->fd);
  299.     sa->fd = -1;
  300.     printf("main: waiting for disk writer\n");
  301.     pthread_join(dwt_id, NULL);
  302.  
  303. qrecord_done:
  304.     free(hdr);
  305.     return 0;
  306. }
  307.  
  308. void *disk_flush_thread(void *argument)
  309. {
  310.     struct qrecord_data *qr = (struct qrecord_data *)argument;
  311.     //fdatasync(qr->filehandle);
  312.     //posix_fadvise(qr->filehandle, 0, 0, POSIX_FADV_NOREUSE);
  313.     posix_fadvise(qr->filehandle, 0, 0, POSIX_FADV_DONTNEED);
  314.  
  315.     pthread_exit(NULL);
  316.     return NULL;
  317. }
  318.  
  319. void *disk_write_thread(void *argument)
  320. {
  321.     struct qrecord_data *qr = (struct qrecord_data *)argument;
  322.     struct sstrx_acqd_data_header *output_buffer = malloc(qr->buffer_size);
  323.     int rv = 0;
  324.     int pac = 0, max_queue_size = 0;
  325.     char filename[MAX_FILENAME + MAX_PATHNAME];
  326.     double old_time = NAN, cur_time;
  327.     struct timeval tv;
  328.     uint64_t old_bytes_written = 0;
  329.     pthread_t dft_id = 0;
  330.  
  331.     while (1) {
  332.         if (0 > queue_try_remove_timeout(&qr->queue, output_buffer, qr->buffer_size, 1000)) {
  333.             if (g_terminate_flag) break;
  334.             continue;
  335.         }
  336.  
  337.         switch (output_buffer->id) {
  338.         case SSTRX_ID_DATA:
  339.             rv = process_data((struct sstrx_acqd_data_header *)output_buffer, qr);
  340.             break;
  341.         case SSTRX_RADAR_ID_EVENT:
  342.             if (qr->automatic_mode) {
  343.                 struct sstrx_radar_event *event = (struct sstrx_radar_event *)output_buffer;
  344.                 if ((event->flags & SSTRX_RADAR_EVENT_START_TASK) && (-1 == qr->filehandle)) {
  345.                     time_t curtime = time(NULL);
  346.                     strftime(qr->filename, MAX_FILENAME, "TS_%Y%m%d_%H%M%S.ts", gmtime(&curtime));
  347.                     stpncpy(filename, qr->output_path, MAX_FILENAME);
  348.                     strncat(filename, qr->filename, MAX_FILENAME);
  349.                     int flags = O_WRONLY | O_CREAT | O_LARGEFILE | O_TRUNC;
  350.                     mode_t mode = S_IWUSR | S_IRUSR | S_IRGRP | S_IROTH;
  351.  
  352.                     qr->filehandle = open(filename, flags, mode);
  353.                     if (-1 == qr->filehandle) {
  354.                         fprintf(stderr, "Can't create file %s: %s\n", qr->filename, strerror(errno));
  355.                     }
  356.                     else {
  357.                         qr->pulses_dropped = 0;
  358.                         pac = 0;
  359.                         qr->bytes_written = 0;
  360.                         qr->bytes_since_flush = 0;
  361.                         printf("\n");
  362.                     }
  363.                     rv = process_header((union sstrx_acqd_cmd *)output_buffer, qr);
  364.                 }
  365.                 else if ((event->flags & SSTRX_RADAR_EVENT_END_TASK) && (-1 != qr->filehandle)) {
  366.                     rv = process_header((union sstrx_acqd_cmd *)output_buffer, qr);
  367.                     strcpy(qr->filename, "<none>");
  368.                     close(qr->filehandle);
  369.                     qr->filehandle = -1;
  370.                     pac = 0;
  371.                     printf("\n");
  372.                 }
  373.                 else {
  374.                     rv = process_header((union sstrx_acqd_cmd *)output_buffer, qr);
  375.                 }
  376.             }
  377.             else {
  378.                 rv = process_header((union sstrx_acqd_cmd *)output_buffer, qr);
  379.             }
  380.             break;
  381.         default:
  382.             rv = process_header((union sstrx_acqd_cmd *)output_buffer, qr);
  383.             break;
  384.         }
  385.         if (0 > rv) {
  386.             fprintf(stderr, "dwt: Couldn't write data: %s\n", strerror(errno));
  387.             break;
  388.         }
  389.         if (g_terminate_flag && queue_get_empty(&qr->queue)) {
  390.             break;
  391.         }
  392.  
  393.         pac--;
  394.         if (max_queue_size < queue_get_size(&qr->queue)) {
  395.             max_queue_size = queue_get_size(&qr->queue);
  396.         }
  397.         if (pac <= 0) {
  398.             putchar('[');
  399.             int bar_size = 20, ctr;
  400.             for (ctr = 0; ctr < (max_queue_size * bar_size) / qr->queue_size; ctr++) {
  401.                 if (qr->queue_full) {
  402.                     putchar('F');
  403.                 }
  404.                 else {
  405.                     putchar('*');
  406.                 }
  407.             }
  408.             for (; ctr < bar_size; ctr++) {
  409.                 putchar(' ');
  410.             }
  411.             printf("] %04d dropped %d ", max_queue_size, qr->pulses_dropped);
  412.             if (qr->filehandle != -1) {
  413.                 gettimeofday(&tv, NULL);
  414.                 cur_time = tv.tv_sec + tv.tv_usec * 1e-6;
  415.                 uint64_t bytes = qr->bytes_written - old_bytes_written;
  416.                 double data_rate = (bytes / 1048576.0) / (cur_time - old_time);
  417.                 printf("%7.1f MB written to %s at %6.2f MB/s\r",
  418.                         qr->bytes_written / 1048576.0, qr->filename,
  419.                         data_rate);
  420.                 old_time = cur_time;
  421.                 old_bytes_written = qr->bytes_written;
  422.             }
  423.             else {
  424.                 printf(" IDLE\r");
  425.             }
  426.             fflush(stdout);
  427.             pac = 64;
  428.             max_queue_size = 0;
  429.         }
  430.         if (qr->bytes_since_flush > qr->flush_after_bytes) {
  431.             if (qr->filehandle != -1) {
  432.                 if (dft_id != 0) {
  433.                     pthread_join(dft_id, NULL);
  434.                 }
  435.                 qr->bytes_since_flush = 0;
  436.                 pthread_create(&dft_id, NULL, disk_flush_thread, qr);
  437.             }
  438.         }
  439.     }
  440.  
  441.     if (-1 == qr->filehandle) {
  442.         close(qr->filehandle);
  443.         qr->filehandle = -1;
  444.     }
  445.     return NULL;
  446. }
  447.  
  448. static int process_data(struct sstrx_acqd_data_header *hdr, struct qrecord_data *qr)
  449. {
  450.     //int ctr;
  451.     //struct sstrx_acqd_data_24b_2ch_sample *data = (struct sstrx_acqd_data_24b_2ch_sample *)(hdr + 1);
  452.  
  453.     static int last_gate_number;
  454.     static int old_clock_offset_ns;
  455.     //static int old_timestamp;
  456.  
  457.     if (qr->show_timestamps) {
  458.         printf("Recording timestamp: %10lld, offset: %8d, delta: %6d\n",
  459.                 (long long int)hdr->timestamp_seconds,
  460.                 hdr->clock_offset_ns, hdr->clock_offset_ns - old_clock_offset_ns);
  461.     }
  462.     old_clock_offset_ns = hdr->clock_offset_ns;
  463.  
  464.     if (((hdr->gate_number - last_gate_number) != 1) &&
  465.             (hdr->gate_number != 0) && (last_gate_number != 0) &&
  466.             ((hdr->gate_number - last_gate_number) > 0)) {
  467.         qr->pulses_dropped += (hdr->gate_number - last_gate_number);
  468.     }
  469.     last_gate_number = hdr->gate_number;
  470.  
  471.     if (-1 == qr->filehandle) {
  472.         return 0;
  473.     }
  474.  
  475.     int byteswritten = write(qr->filehandle, hdr, hdr->len);
  476.     if (-1 == byteswritten) {
  477.         perror("write");
  478.         return -1;
  479.     }
  480.     if (byteswritten != hdr->len) {
  481.         fprintf(stderr, "write only wrote %d bytes of a %d packet\n", byteswritten, hdr->len);
  482.         return -1;
  483.     }
  484.     qr->bytes_written += byteswritten;
  485.     qr->bytes_since_flush += byteswritten;
  486.  
  487.     return 0;
  488. }
  489.  
  490. static int process_header(union sstrx_acqd_cmd *cmd, struct qrecord_data *qr)
  491. {
  492.     struct sstrx_acqd_cmd_generic *hdr = &cmd->generic;
  493.    
  494.     if (-1 == qr->filehandle) {
  495.         return 0;
  496.     }
  497.  
  498.     int byteswritten = write(qr->filehandle, hdr, hdr->len);
  499.     if (-1 == byteswritten) {
  500.         perror("write");
  501.         return -1;
  502.     }
  503.     if (byteswritten != hdr->len) {
  504.         fprintf(stderr, "write only wrote %d bytes of a %d packet\n", byteswritten, hdr->len);
  505.         return -1;
  506.     }
  507.     qr->bytes_written += byteswritten;
  508.  
  509.     return 0;
  510. }
  511.  
  512. void break_handler(int signo)
  513. {
  514.     if (signo == SIGINT) {
  515.         g_terminate_flag = 1;
  516.     }
  517. }
  518.  
  519. static void setup_signals(void)
  520. {
  521.     /*
  522.     signal(SIGINT, break_handler);
  523.     signal(SIGPIPE, SIG_IGN);
  524.     */
  525.     static struct sigaction sa;
  526.  
  527.     sa.sa_handler = break_handler;
  528.     sigemptyset(&sa.sa_mask);
  529.     sa.sa_flags = SA_NOCLDSTOP | SA_RESTART;
  530.  
  531.     if(-1 == sigaction(SIGINT, &sa, NULL)) perror("sigaction");
  532.  
  533.     sa.sa_handler = SIG_IGN;
  534.     sigemptyset(&sa.sa_mask);
  535.     sa.sa_flags = 0;
  536.  
  537.     if(-1 == sigaction(SIGPIPE, &sa, NULL)) perror("sigaction");
  538. }
  539.  
Advertisement
Add Comment
Please, Sign In to add comment