diff --git a/src/main/main.c b/src/main/main.c index d662f8c..7429f60 100644 --- a/src/main/main.c +++ b/src/main/main.c @@ -274,16 +274,26 @@ static int pv__store_and_forward(pvstate_t state, opts_t opts, bool can_have_eta * Run the main transfer loop for the given side of monitor mode, using the * "command_fd" file descriptor as the transfer loop's output (on the "in" * side, it's a pipe to the command) or as its input (on the "out" side, - * it's a pipe from the command). If both sides are active, - * "other_side_pid" is the PID of the other side. + * it's a pipe from the command). + * + * If both sides are active, "othermonitor_pid" is the PID of the monitor on + * the other side, "othermonitor_read_fd" is a pipe file descriptor to read + * info from the other monitor, and "othermonitor_write_fd" is for writing + * info to the other monitor. * * Returns the appropriate exit status. * - * As side effect, "command_fd" is closed. + * As a side effect, "command_fd" is closed. */ static int pv__run_monitor(const char *program_name, pvstate_t state, pvside_t side, int command_fd, /*@unused@ */ - __attribute__((unused)) pid_t other_side_pid) + __attribute__((unused)) pid_t othermonitor_pid, + /*@unused@ */ + __attribute__((unused)) + int othermonitor_read_fd, + /*@unused@ */ + __attribute__((unused)) + int othermonitor_write_fd) { const char *dummy_argv[1]; /* flawfinder: ignore */ @@ -315,9 +325,7 @@ static int pv__run_monitor(const char *program_name, pvstate_t state, pvside_t s fprintf(stderr, "%s: %s\n", program_name, strerror(errno)); } - /* TODO: use other_side_pid, for ratio display. */ - /* TODO: exchange transfer information like -Q. */ - /* TODO: ...or arrange pipes between the two PIDs. */ + /* TODO: pass othermonitor_* into the state for ratio display. */ /*@-observertrans@ */ dummy_argv[0] = "-"; @@ -329,13 +337,21 @@ static int pv__run_monitor(const char *program_name, pvstate_t state, pvside_t s /* - * Run in monitor mode: run a process and run the main transfer loop on its - * input, output, or both. Returns the appropriate exit status. + * Monitor mode: run a process and run the main transfer loop on its input, + * output, or both. Returns the appropriate exit status. + * + * When monitoring both sides, two pipes are set up between the monitor + * processes for each side - in to out, and out to in - for them to exchange + * information about the transfer, so both sides can see how many bytes the + * other side has transferred. This is what allows the in:out ratio to be + * displayed. */ static int pv__monitor(pvstate_t state, opts_t opts) { - int pipefd_in[2]; - int pipefd_out[2]; + int pipefd_cmd_in[2]; /* pipe from the monitor to the command */ + int pipefd_cmd_out[2]; /* pipe from the command to the monitor */ + int pipefd_in_to_out[2]; /* from in monitor to out monitor */ + int pipefd_out_to_in[2]; /* from out monitor to in monitor */ int retcode, pid_status; pid_t command_pid, in_monitor_pid, out_monitor_pid, waited_pid; @@ -353,82 +369,79 @@ static int pv__monitor(pvstate_t state, opts_t opts) * Create the pipes for communicating with the command. */ - pipefd_in[0] = -1; - pipefd_in[1] = -1; - pipefd_out[0] = -1; - pipefd_out[1] = -1; + pipefd_cmd_in[0] = -1; + pipefd_cmd_in[1] = -1; + pipefd_cmd_out[0] = -1; + pipefd_cmd_out[1] = -1; /* Pipe for the input side of the command, if we're monitoring it. */ if ((PV_SIDE_IN == opts->side) || (PV_SIDE_BOTH == opts->side)) { - if (0 != pipe(pipefd_in)) { + if (0 != pipe(pipefd_cmd_in)) { fprintf(stderr, "%s: %s\n", opts->program_name, strerror(errno)); return PV_ERROREXIT_MONITOR; } - debug("pipefd_in[]=(%d,%d)", pipefd_in[0], pipefd_in[1]); + debug("pipefd_cmd_in[]=(%d,%d)", pipefd_cmd_in[0], pipefd_cmd_in[1]); } /* Pipe for the output side of the command, if we're monitoring it. */ if ((PV_SIDE_OUT == opts->side) || (PV_SIDE_BOTH == opts->side)) { - if (0 != pipe(pipefd_out)) { + if (0 != pipe(pipefd_cmd_out)) { fprintf(stderr, "%s: %s\n", opts->program_name, strerror(errno)); - if (-1 != pipefd_in[0]) - (void) close(pipefd_in[0]); - if (-1 != pipefd_in[1]) - (void) close(pipefd_in[1]); + if (-1 != pipefd_cmd_in[0]) + (void) close(pipefd_cmd_in[0]); + if (-1 != pipefd_cmd_in[1]) + (void) close(pipefd_cmd_in[1]); return PV_ERROREXIT_MONITOR; } - debug("pipefd_out[]=(%d,%d)", pipefd_out[0], pipefd_out[1]); + debug("pipefd_cmd_out[]=(%d,%d)", pipefd_cmd_out[0], pipefd_cmd_out[1]); } + /* Common idiom to close an fd, and set it to -1, if it's open. */ +#define close_if_open(x) if (-1 != x) { \ +(void) close(x); \ +x = 1; \ +} + /* Create a process to run the command. */ command_pid = (pid_t) fork(); if (command_pid < 0) { fprintf(stderr, "%s: %s\n", opts->program_name, strerror(errno)); - if (-1 != pipefd_in[0]) - (void) close(pipefd_in[0]); - if (-1 != pipefd_in[1]) - (void) close(pipefd_in[1]); - if (-1 != pipefd_out[0]) - (void) close(pipefd_out[0]); - if (-1 != pipefd_out[1]) - (void) close(pipefd_out[1]); + close_if_open(pipefd_cmd_in[0]); + close_if_open(pipefd_cmd_in[1]); + close_if_open(pipefd_cmd_out[0]); + close_if_open(pipefd_cmd_out[1]); return PV_ERROREXIT_MONITOR; } else if (0 == command_pid) { /* Command process. */ /* Close the write end of the "in" pipe. */ - if (-1 != pipefd_in[1]) { - (void) close(pipefd_in[1]); - pipefd_in[1] = -1; - } + close_if_open(pipefd_cmd_in[1]); + /* Close the read end of the "out" pipe. */ - if (-1 != pipefd_out[0]) { - (void) close(pipefd_out[0]); - pipefd_out[0] = -1; - } + close_if_open(pipefd_cmd_out[0]); /* Put the read end of the "in" pipe on stdin. */ - if (-1 != pipefd_in[0]) { - debug("replacing command stdin with fd %d", pipefd_in[0]); - if (dup2(pipefd_in[0], STDIN_FILENO) < 0) { + if (-1 != pipefd_cmd_in[0]) { + debug("replacing command stdin with fd %d", pipefd_cmd_in[0]); + if (dup2(pipefd_cmd_in[0], STDIN_FILENO) < 0) { fprintf(stderr, "%s: %s\n", opts->program_name, strerror(errno)); exit(EXIT_FAILURE); } - (void) close(pipefd_in[0]); - debug("replaced command stdin with fd %d", pipefd_in[0]); - pipefd_in[0] = -1; + (void) close(pipefd_cmd_in[0]); + debug("replaced command stdin with fd %d", pipefd_cmd_in[0]); + pipefd_cmd_in[0] = -1; } /* Put the write end of the "out" pipe on stdout. */ - if (-1 != pipefd_out[1]) { - debug("replacing command stdout with fd %d", pipefd_out[1]); - if (dup2(pipefd_out[1], STDOUT_FILENO) < 0) { + if (-1 != pipefd_cmd_out[1]) { + debug("replacing command stdout with fd %d", pipefd_cmd_out[1]); + if (dup2(pipefd_cmd_out[1], STDOUT_FILENO) < 0) { fprintf(stderr, "%s: %s\n", opts->program_name, strerror(errno)); exit(EXIT_FAILURE); } - (void) close(pipefd_out[1]); - debug("replaced command stdout with fd %d", pipefd_out[1]); - pipefd_out[1] = -1; + (void) close(pipefd_cmd_out[1]); + debug("replaced command stdout with fd %d", pipefd_cmd_out[1]); + pipefd_cmd_out[1] = -1; } /* Execute the command. */ @@ -447,14 +460,38 @@ static int pv__monitor(pvstate_t state, opts_t opts) /* Main process, not the command process. */ /* Close the read end of the "in" pipe. */ - if (-1 != pipefd_in[0]) { - (void) close(pipefd_in[0]); - pipefd_in[0] = -1; - } + close_if_open(pipefd_cmd_in[0]); + /* Close the write end of the "out" pipe. */ - if (-1 != pipefd_out[1]) { - (void) close(pipefd_out[1]); - pipefd_out[1] = -1; + close_if_open(pipefd_cmd_out[1]); + + /* + * If monitoring both input and output, create pipes for the two + * monitors to talk to each other in both directions. + */ + pipefd_in_to_out[0] = -1; + pipefd_in_to_out[1] = -1; + pipefd_out_to_in[0] = -1; + pipefd_out_to_in[1] = -1; + if (PV_SIDE_BOTH == opts->side) { + if (0 != pipe(pipefd_in_to_out)) { + fprintf(stderr, "%s: %s\n", opts->program_name, strerror(errno)); + close_if_open(pipefd_cmd_in[1]); + close_if_open(pipefd_cmd_out[0]); + (void) kill(command_pid, SIGTERM); + return PV_ERROREXIT_MONITOR; + } + debug("pipefd_in_to_out[]=(%d,%d)", pipefd_in_to_out[0], pipefd_in_to_out[1]); + if (0 != pipe(pipefd_out_to_in)) { + fprintf(stderr, "%s: %s\n", opts->program_name, strerror(errno)); + close_if_open(pipefd_in_to_out[0]); + close_if_open(pipefd_in_to_out[1]); + close_if_open(pipefd_cmd_in[1]); + close_if_open(pipefd_cmd_out[0]); + (void) kill(command_pid, SIGTERM); + return PV_ERROREXIT_MONITOR; + } + debug("pipefd_out_to_in[]=(%d,%d)", pipefd_out_to_in[0], pipefd_out_to_in[1]); } /* @@ -468,21 +505,44 @@ static int pv__monitor(pvstate_t state, opts_t opts) out_monitor_pid = (pid_t) fork(); if (out_monitor_pid < 0) { fprintf(stderr, "%s: %s\n", opts->program_name, strerror(errno)); - if (-1 != pipefd_in[1]) - (void) close(pipefd_in[1]); - if (-1 != pipefd_out[0]) - (void) close(pipefd_out[0]); + close_if_open(pipefd_in_to_out[0]); + close_if_open(pipefd_in_to_out[1]); + close_if_open(pipefd_out_to_in[0]); + close_if_open(pipefd_out_to_in[1]); + close_if_open(pipefd_cmd_in[1]); + close_if_open(pipefd_cmd_out[0]); (void) kill(command_pid, SIGTERM); return PV_ERROREXIT_MONITOR; } else if (0 == out_monitor_pid) { /* Output side monitoring process. */ - /* Close the write end of the "in" pipe. */ - if (-1 != pipefd_in[1]) - (void) close(pipefd_in[1]); + /* Close the write end of the command "in" pipe. */ + close_if_open(pipefd_cmd_in[1]); - return pv__run_monitor(opts->program_name, state, PV_SIDE_OUT, pipefd_out[0], in_monitor_pid); + /* Close the write end of the in-to-out pipe. */ + close_if_open(pipefd_in_to_out[1]); + + /* Close the read end of the out-to-in pipe. */ + close_if_open(pipefd_out_to_in[0]); + + retcode = + pv__run_monitor(opts->program_name, state, PV_SIDE_OUT, pipefd_cmd_out[0], in_monitor_pid, + pipefd_in_to_out[0], pipefd_out_to_in[1]); + + /* Close the other ends of the intra-monitor pipes. */ + close_if_open(pipefd_in_to_out[0]); + close_if_open(pipefd_out_to_in[1]); + + return retcode; } + + /* Input side monitoring process. */ + + /* Close the read end of the in-to-out pipe. */ + close_if_open(pipefd_in_to_out[0]); + + /* Close the write end of the out-to-in pipe. */ + close_if_open(pipefd_out_to_in[1]); } /* Monitor the remaining side. */ @@ -506,17 +566,17 @@ static int pv__monitor(pvstate_t state, opts_t opts) #endif case PV_SIDE_IN: /* Close the read end of the "out" pipe. */ - if (-1 != pipefd_out[0]) - (void) close(pipefd_out[0]); - pipefd_out[0] = -1; - retcode = pv__run_monitor(opts->program_name, state, PV_SIDE_IN, pipefd_in[1], out_monitor_pid); + close_if_open(pipefd_cmd_out[0]); + retcode = + pv__run_monitor(opts->program_name, state, PV_SIDE_IN, pipefd_cmd_in[1], out_monitor_pid, + pipefd_out_to_in[0], pipefd_in_to_out[1]); break; case PV_SIDE_OUT: /* Close the write end of the "in" pipe. */ - if (-1 != pipefd_in[1]) - (void) close(pipefd_in[1]); - pipefd_in[1] = -1; - retcode = pv__run_monitor(opts->program_name, state, PV_SIDE_OUT, pipefd_out[0], in_monitor_pid); + close_if_open(pipefd_cmd_in[1]); + retcode = + pv__run_monitor(opts->program_name, state, PV_SIDE_OUT, pipefd_cmd_out[0], in_monitor_pid, + pipefd_in_to_out[0], pipefd_out_to_in[1]); break; } @@ -542,6 +602,10 @@ static int pv__monitor(pvstate_t state, opts_t opts) /* If monitoring both sides, wait for the "out" side to exit. */ if (PV_SIDE_BOTH == opts->side && -1 != out_monitor_pid) { + /* Close our ends of the intra-monitor pipes. */ + close_if_open(pipefd_in_to_out[1]); + close_if_open(pipefd_out_to_in[0]); + /* Wait for the "out" side to exit. */ do { /*@-type@ */ /* splint disagreement about __pid_t vs pid_t. */