Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- /**
- * \file src/sstrx_qrecord.c
- * \brief Main module for quick timeseries recorder
- * \author *
- *
- * Implements a quick timeseries data recorder client for sstrx
- *
- */
- #include "config.h"
- #include <stdio.h>
- #include <math.h>
- #include <syslog.h>
- #include <stdlib.h>
- #include <string.h>
- #include <unistd.h>
- #include <getopt.h>
- #include <pthread.h>
- #include <signal.h>
- #include <sys/socket.h>
- #include <sys/types.h>
- #include <sys/stat.h>
- #include <fcntl.h>
- #include <sys/types.h>
- #include <sys/time.h>
- #include <errno.h>
- #include <time.h>
- #include <poll.h>
- #include <limits.h>
- #include "net_utils.h"
- #include "client_lib.h"
- #include "data_format.h"
- #include "radar_info.h"
- #include "queue.h"
- /* * * * Module-local constants * * * */
- enum PROGRAM_ARGS {
- ARG_HELP = 256,
- ARG_SERVERNAME,
- ARG_TIMEOUT,
- ARG_SAMPLES,
- ARG_QUEUESIZE,
- ARG_NUM_PRTS,
- ARG_OUTPUT_DIR,
- ARG_FLUSHAFTER,
- };
- #define MAX_HOSTNAME 128
- #define MAX_FILENAME PATH_MAX
- #define MAX_PATHNAME PATH_MAX
- #define DEFAULT_QUEUE_SIZE 2000
- #define DEFAULT_SAMPLES 5000
- #define DEFAULT_FLUSHAFTER_MB 32
- /* * * * Module-local structure definitions * * * */
- struct qrecord_data {
- uint32_t num_samples;
- char filename[MAX_FILENAME];
- int filehandle;
- struct queue queue;
- int buffer_size;
- int queue_size;
- int num_prts;
- int show_timestamps;
- uint32_t automatic_mode;
- uint64_t bytes_written;
- char output_path[MAX_PATHNAME];
- int queue_full;
- int pulses_dropped;
- uint64_t flush_after_bytes;
- uint64_t bytes_since_flush;
- };
- /* * * * Module-local function declarations * * * */
- static int do_qrecord(struct sstrx_acqd *sa, struct qrecord_data *qr);
- static void setup_signals(void);
- static int process_header(union sstrx_acqd_cmd *hdr, struct qrecord_data *qr);
- static int process_data(struct sstrx_acqd_data_header *hdr, struct qrecord_data *qr);
- void *disk_write_thread(void *argument);
- /* * * * Module-local variables * * * */
- static char hostname[MAX_HOSTNAME] = "localhost";
- static uint32_t port = SSTRX_ACQD_DEFAULT_PORT;
- volatile static int g_terminate_flag = 0;
- static struct qrecord_data qrecord_data = {
- .num_samples = DEFAULT_SAMPLES,
- .num_prts = 0,
- .queue_size = DEFAULT_QUEUE_SIZE,
- .flush_after_bytes = 1048576 * DEFAULT_FLUSHAFTER_MB,
- };
- /** Usage string, printed when user needs help w/ command line */
- static char *usage =
- "SSTRX quick timeseries recorder\n"
- "sstrx-acqd-qrecord [options] [filename]\n"
- " [filename]: File to record to.\n"
- " If no file is specified, enter automatic mode where the recorder\n"
- " waits for event notifications to start and stop recording, files\n"
- " are created with the names TS_yyyymmdd_hhmmss.dat"
- "--help: Show usage information\n"
- "--server=<hostname:port>: Set the server name and port. Defaults to\n"
- " env. var. ACQD_SVR if available, else localhost\n"
- "--timeout=<timeout>: Set the timeout in seconds (fractional OK)\n"
- "--samples=<samples>: Acquire <samples> samples per PRT, defaults to " xstr(DEFAULT_SAMPLES) "\n"
- "--queuesize=<size>: Set queue size to <size>, defaults to " xstr(DEFAULT_QUEUE_SIZE) "\n"
- "--prts=<prts>: Stop recording after <prts> PRTs.\n"
- "--showtimestamps: Display timestamp information\n"
- "--outdir=<path>: Write to specified directory (default: current dir)\n"
- "--flushafter=<mbytes>: Flush cache after specified megabytes, default " xstr(DEFAULT_FLUSHAFTER_MB) "M \n"
- "\n";
- /* * * * Function definitions * * * */
- int main(int argc, char *argv[])
- {
- int opt, opt_idx, opt_errors = 0;
- int rv = -1;
- struct option opt_lst[] = {
- {"help", no_argument, 0, ARG_HELP},
- {"server", required_argument, 0, ARG_SERVERNAME},
- {"timeout", required_argument, 0, ARG_TIMEOUT},
- {"samples", required_argument, 0, ARG_SAMPLES},
- {"queuesize", required_argument, 0, ARG_QUEUESIZE},
- {"prts", required_argument, 0, ARG_NUM_PRTS},
- {"showtimestamps", no_argument, &qrecord_data.show_timestamps, 1},
- {"outdir", required_argument, 0, ARG_OUTPUT_DIR},
- {"flushafter", required_argument, 0, ARG_FLUSHAFTER},
- {0, 0, 0, 0} // Sentinel
- };
- struct sstrx_acqd sstrx_acqd;
- struct sstrx_acqd *sa = &sstrx_acqd;
- if (NULL != getenv("ACQD_SVR")) {
- strncpy(hostname, getenv("ACQD_SVR"), MAX_HOSTNAME-1);
- strtok(hostname, ":");
- char *port_str = strtok(NULL, ":");
- port = (port_str != 0) ? strtoul(port_str, NULL, 10) : SSTRX_ACQD_DEFAULT_PORT;
- }
- sa->timeout = 1000;
- sa->max_skip = 100;
- sa->input_cmd_buffer_size = sstrx_acqd_buffer_size(SSTRX_ACQD_BUFSZ_ALL);
- //printf("Allocating buffer of size %d bytes\n", sa->input_cmd_buffer_size);
- sa->input_cmd_buffer = malloc(sa->input_cmd_buffer_size);
- struct qrecord_data *qr = &qrecord_data;
- while (-1 != (opt = getopt_long(argc, argv, "", opt_lst, &opt_idx))) {
- switch (opt) {
- case ARG_HELP:
- puts(usage);
- return 0;
- break;
- case ARG_SERVERNAME:
- {
- strncpy(hostname, strtok(optarg, ":"), MAX_HOSTNAME-1);
- char *port_str = strtok(NULL, ":");
- port = (port_str != 0) ? strtoul(port_str, NULL, 10) : SSTRX_ACQD_DEFAULT_PORT;
- }
- break;
- case ARG_TIMEOUT:
- sa->timeout = round(strtod(optarg, NULL) * 1e3);
- break;
- case ARG_SAMPLES:
- qr->num_samples = strtoul(optarg, NULL, 10);
- break;
- case ARG_QUEUESIZE:
- qr->queue_size = strtoul(optarg, NULL, 10);
- break;
- case ARG_NUM_PRTS:
- qr->num_prts = strtoul(optarg, NULL, 10);
- break;
- case ARG_OUTPUT_DIR:
- stpncpy(qr->output_path, optarg, MAX_PATHNAME-1);
- strncat(qr->output_path, "/", MAX_PATHNAME);
- break;
- case ARG_FLUSHAFTER:
- qr->flush_after_bytes = strtoul(optarg, NULL, 10) * 1048576;
- break;
- case 0:
- break;
- default:
- opt_errors++;
- break;
- }
- }
- qr->buffer_size = sa->input_cmd_buffer_size;
- if (0 != queue_init(&qr->queue, qr->queue_size, sa->input_cmd_buffer_size)) {
- fprintf(stderr, "Queue init failed\n");
- goto terminate_with_error;
- }
- if (opt_errors) {
- puts(usage);
- return -1;
- }
- if (optind < argc) {
- qr->automatic_mode = 0;
- strncpy(qr->filename, argv[optind], MAX_FILENAME);
- int flags = O_WRONLY | O_CREAT | O_LARGEFILE | O_TRUNC;
- mode_t mode = S_IWUSR | S_IRUSR | S_IRGRP | S_IROTH;
- qr->filehandle = open(qr->filename, flags, mode);
- printf("open returns: %d\n", qr->filehandle);
- if (qr->filehandle < 0) {
- perror("open");
- }
- }
- else {
- qr->automatic_mode = 1;
- strcpy(qr->filename, "<none>");
- qr->filehandle = -1;
- }
- printf("Connecting to %s:%d\n", hostname, port);
- sa->fd = open_socket(hostname, port);
- if (-1 == sa->fd) {
- fprintf(stderr, "Could not connect to %s:%d\n", hostname, port);
- return -1;
- }
- setup_signals();
- rv = sstrx_acqd_set_data_format(sa,
- SSTRX_ACQD_FLAGS_FMT_24B_2CH,
- qr->num_samples, 0);
- if(0 > rv) goto terminate_with_error;
- rv = do_qrecord(sa, qr);
- if(0 > rv) goto terminate_with_error;
- if (sa->fd != -1) {
- shutdown(sa->fd, SHUT_RDWR);
- close(sa->fd);
- }
- return 0;
- terminate_with_error:
- if (sa->fd != -1) {
- close(sa->fd);
- }
- if (0 == rv) {
- printf("No connection to server %s:%d\n", hostname, port);
- }
- else {
- printf("Error: %s\n", enum_to_string(sstrx_acqd_client_errors, rv));
- }
- return rv;
- }
- static int do_qrecord(struct sstrx_acqd *sa, struct qrecord_data *qr)
- {
- struct sstrx_acqd_data_header *hdr = (struct sstrx_acqd_data_header *)sa->input_cmd_buffer;
- int prts_acquired = 0;
- pthread_t dwt_id;
- pthread_create(&dwt_id, NULL, disk_write_thread, qr);
- int bytesread, byteswanted;
- while(!g_terminate_flag) {
- struct pollfd pfd = {sa->fd, POLLIN, 0};
- int rv = poll(&pfd, 1, 1000);
- if (1 != rv) {
- if (g_terminate_flag) break;
- continue;
- }
- byteswanted = sizeof(struct sstrx_acqd_cmd_generic);
- bytesread = safe_read(sa->fd, hdr, byteswanted);
- if (bytesread != byteswanted) {
- fprintf(stderr, "error reading ID/len: %d\n", bytesread);
- goto qrecord_done;
- }
- void *write_location = ((struct sstrx_acqd_cmd_generic *)hdr) + 1;
- byteswanted = hdr->len - sizeof(struct sstrx_acqd_cmd_generic);
- uint32_t temp;
- bytesread = safe_read_max(sa->fd, write_location, byteswanted,
- sa->input_cmd_buffer_size, &temp, sizeof(temp));
- if (bytesread < 0) {
- perror("read");
- goto qrecord_done;
- }
- if (0 != queue_try_add(&qr->queue, sa->input_cmd_buffer, hdr->len)) {
- qr->queue_full = 1;
- }
- else {
- qr->queue_full = 0;
- }
- prts_acquired++;
- if (qr-> num_prts && (prts_acquired >= qr->num_prts)) {
- g_terminate_flag = 1;
- break;
- }
- }
- close(sa->fd);
- sa->fd = -1;
- printf("main: waiting for disk writer\n");
- pthread_join(dwt_id, NULL);
- qrecord_done:
- free(hdr);
- return 0;
- }
- void *disk_flush_thread(void *argument)
- {
- struct qrecord_data *qr = (struct qrecord_data *)argument;
- //fdatasync(qr->filehandle);
- //posix_fadvise(qr->filehandle, 0, 0, POSIX_FADV_NOREUSE);
- posix_fadvise(qr->filehandle, 0, 0, POSIX_FADV_DONTNEED);
- pthread_exit(NULL);
- return NULL;
- }
- void *disk_write_thread(void *argument)
- {
- struct qrecord_data *qr = (struct qrecord_data *)argument;
- struct sstrx_acqd_data_header *output_buffer = malloc(qr->buffer_size);
- int rv = 0;
- int pac = 0, max_queue_size = 0;
- char filename[MAX_FILENAME + MAX_PATHNAME];
- double old_time = NAN, cur_time;
- struct timeval tv;
- uint64_t old_bytes_written = 0;
- pthread_t dft_id = 0;
- while (1) {
- if (0 > queue_try_remove_timeout(&qr->queue, output_buffer, qr->buffer_size, 1000)) {
- if (g_terminate_flag) break;
- continue;
- }
- switch (output_buffer->id) {
- case SSTRX_ID_DATA:
- rv = process_data((struct sstrx_acqd_data_header *)output_buffer, qr);
- break;
- case SSTRX_RADAR_ID_EVENT:
- if (qr->automatic_mode) {
- struct sstrx_radar_event *event = (struct sstrx_radar_event *)output_buffer;
- if ((event->flags & SSTRX_RADAR_EVENT_START_TASK) && (-1 == qr->filehandle)) {
- time_t curtime = time(NULL);
- strftime(qr->filename, MAX_FILENAME, "TS_%Y%m%d_%H%M%S.ts", gmtime(&curtime));
- stpncpy(filename, qr->output_path, MAX_FILENAME);
- strncat(filename, qr->filename, MAX_FILENAME);
- int flags = O_WRONLY | O_CREAT | O_LARGEFILE | O_TRUNC;
- mode_t mode = S_IWUSR | S_IRUSR | S_IRGRP | S_IROTH;
- qr->filehandle = open(filename, flags, mode);
- if (-1 == qr->filehandle) {
- fprintf(stderr, "Can't create file %s: %s\n", qr->filename, strerror(errno));
- }
- else {
- qr->pulses_dropped = 0;
- pac = 0;
- qr->bytes_written = 0;
- qr->bytes_since_flush = 0;
- printf("\n");
- }
- rv = process_header((union sstrx_acqd_cmd *)output_buffer, qr);
- }
- else if ((event->flags & SSTRX_RADAR_EVENT_END_TASK) && (-1 != qr->filehandle)) {
- rv = process_header((union sstrx_acqd_cmd *)output_buffer, qr);
- strcpy(qr->filename, "<none>");
- close(qr->filehandle);
- qr->filehandle = -1;
- pac = 0;
- printf("\n");
- }
- else {
- rv = process_header((union sstrx_acqd_cmd *)output_buffer, qr);
- }
- }
- else {
- rv = process_header((union sstrx_acqd_cmd *)output_buffer, qr);
- }
- break;
- default:
- rv = process_header((union sstrx_acqd_cmd *)output_buffer, qr);
- break;
- }
- if (0 > rv) {
- fprintf(stderr, "dwt: Couldn't write data: %s\n", strerror(errno));
- break;
- }
- if (g_terminate_flag && queue_get_empty(&qr->queue)) {
- break;
- }
- pac--;
- if (max_queue_size < queue_get_size(&qr->queue)) {
- max_queue_size = queue_get_size(&qr->queue);
- }
- if (pac <= 0) {
- putchar('[');
- int bar_size = 20, ctr;
- for (ctr = 0; ctr < (max_queue_size * bar_size) / qr->queue_size; ctr++) {
- if (qr->queue_full) {
- putchar('F');
- }
- else {
- putchar('*');
- }
- }
- for (; ctr < bar_size; ctr++) {
- putchar(' ');
- }
- printf("] %04d dropped %d ", max_queue_size, qr->pulses_dropped);
- if (qr->filehandle != -1) {
- gettimeofday(&tv, NULL);
- cur_time = tv.tv_sec + tv.tv_usec * 1e-6;
- uint64_t bytes = qr->bytes_written - old_bytes_written;
- double data_rate = (bytes / 1048576.0) / (cur_time - old_time);
- printf("%7.1f MB written to %s at %6.2f MB/s\r",
- qr->bytes_written / 1048576.0, qr->filename,
- data_rate);
- old_time = cur_time;
- old_bytes_written = qr->bytes_written;
- }
- else {
- printf(" IDLE\r");
- }
- fflush(stdout);
- pac = 64;
- max_queue_size = 0;
- }
- if (qr->bytes_since_flush > qr->flush_after_bytes) {
- if (qr->filehandle != -1) {
- if (dft_id != 0) {
- pthread_join(dft_id, NULL);
- }
- qr->bytes_since_flush = 0;
- pthread_create(&dft_id, NULL, disk_flush_thread, qr);
- }
- }
- }
- if (-1 == qr->filehandle) {
- close(qr->filehandle);
- qr->filehandle = -1;
- }
- return NULL;
- }
- static int process_data(struct sstrx_acqd_data_header *hdr, struct qrecord_data *qr)
- {
- //int ctr;
- //struct sstrx_acqd_data_24b_2ch_sample *data = (struct sstrx_acqd_data_24b_2ch_sample *)(hdr + 1);
- static int last_gate_number;
- static int old_clock_offset_ns;
- //static int old_timestamp;
- if (qr->show_timestamps) {
- printf("Recording timestamp: %10lld, offset: %8d, delta: %6d\n",
- (long long int)hdr->timestamp_seconds,
- hdr->clock_offset_ns, hdr->clock_offset_ns - old_clock_offset_ns);
- }
- old_clock_offset_ns = hdr->clock_offset_ns;
- if (((hdr->gate_number - last_gate_number) != 1) &&
- (hdr->gate_number != 0) && (last_gate_number != 0) &&
- ((hdr->gate_number - last_gate_number) > 0)) {
- qr->pulses_dropped += (hdr->gate_number - last_gate_number);
- }
- last_gate_number = hdr->gate_number;
- if (-1 == qr->filehandle) {
- return 0;
- }
- int byteswritten = write(qr->filehandle, hdr, hdr->len);
- if (-1 == byteswritten) {
- perror("write");
- return -1;
- }
- if (byteswritten != hdr->len) {
- fprintf(stderr, "write only wrote %d bytes of a %d packet\n", byteswritten, hdr->len);
- return -1;
- }
- qr->bytes_written += byteswritten;
- qr->bytes_since_flush += byteswritten;
- return 0;
- }
- static int process_header(union sstrx_acqd_cmd *cmd, struct qrecord_data *qr)
- {
- struct sstrx_acqd_cmd_generic *hdr = &cmd->generic;
- if (-1 == qr->filehandle) {
- return 0;
- }
- int byteswritten = write(qr->filehandle, hdr, hdr->len);
- if (-1 == byteswritten) {
- perror("write");
- return -1;
- }
- if (byteswritten != hdr->len) {
- fprintf(stderr, "write only wrote %d bytes of a %d packet\n", byteswritten, hdr->len);
- return -1;
- }
- qr->bytes_written += byteswritten;
- return 0;
- }
- void break_handler(int signo)
- {
- if (signo == SIGINT) {
- g_terminate_flag = 1;
- }
- }
- static void setup_signals(void)
- {
- /*
- signal(SIGINT, break_handler);
- signal(SIGPIPE, SIG_IGN);
- */
- static struct sigaction sa;
- sa.sa_handler = break_handler;
- sigemptyset(&sa.sa_mask);
- sa.sa_flags = SA_NOCLDSTOP | SA_RESTART;
- if(-1 == sigaction(SIGINT, &sa, NULL)) perror("sigaction");
- sa.sa_handler = SIG_IGN;
- sigemptyset(&sa.sa_mask);
- sa.sa_flags = 0;
- if(-1 == sigaction(SIGPIPE, &sa, NULL)) perror("sigaction");
- }
Advertisement
Add Comment
Please, Sign In to add comment