1499
{
1500
return child_is_working(pp_child) && pp_child->process.in > 0;
1501
}
1502
+static int child_is_sending_output(const struct parallel_child *pp_child)
1503
+{
1504
+ /*
1505
+ * all pp children which buffer output through run_command via ungroup=0
1506
+ * redirect stdout to stderr, so we just need to check process.err.
1507
+ */
1508
+ return child_is_working(pp_child) && pp_child->process.err > 0;
1509
+}
1510
1511
struct parallel_processes {
1512
size_t nr_processes;
1570
1571
CALLOC_ARRAY(pp->children, n);
1572
if (!opts->ungroup)
1565
- CALLOC_ARRAY(pp->pfd, n);
1573
+ CALLOC_ARRAY(pp->pfd, n * 2);
1574
1575
for (size_t i = 0; i < n; i++) {
1576
strbuf_init(&pp->children[i].err, 0);
1715
}
1716
}
1717
1710
-static void pp_buffer_stderr(struct parallel_processes *pp,
1711
- const struct run_process_parallel_opts *opts,
1712
- int output_timeout)
1718
+static void pp_buffer_io(struct parallel_processes *pp,
1719
+ const struct run_process_parallel_opts *opts,
1720
+ int timeout)
1721
{
1714
- while (poll(pp->pfd, opts->processes, output_timeout) < 0) {
1722
+ /* for each potential child slot, prepare two pollfd entries */
1723
+ for (size_t i = 0; i < opts->processes; i++) {
1724
+ if (child_is_sending_output(&pp->children[i])) {
1725
+ pp->pfd[2*i].fd = pp->children[i].process.err;
1726
+ pp->pfd[2*i].events = POLLIN | POLLHUP;
1727
+ } else {
1728
+ pp->pfd[2*i].fd = -1;
1729
+ }
1730
+
1731
+ if (child_is_receiving_input(&pp->children[i])) {
1732
+ pp->pfd[2*i+1].fd = pp->children[i].process.in;
1733
+ pp->pfd[2*i+1].events = POLLOUT;
1734
+ } else {
1735
+ pp->pfd[2*i+1].fd = -1;
1736
+ }
1737
+ }
1738
+
1739
+ while (poll(pp->pfd, opts->processes * 2, timeout) < 0) {
1740
if (errno == EINTR)
1741
continue;
1742
pp_cleanup(pp, opts);
1743
die_errno("poll");
1744
}
1745
1721
- /* Buffer output from all pipes. */
1746
for (size_t i = 0; i < opts->processes; i++) {
1747
+ /* Handle input feeding (stdin) */
1748
+ if (pp->pfd[2*i+1].revents & (POLLOUT | POLLHUP | POLLERR)) {
1749
+ if (opts->feed_pipe) {
1750
+ int ret = opts->feed_pipe(pp->children[i].process.in,
1751
+ opts->data,
1752
+ pp->children[i].data);
1753
+ if (ret < 0)
1754
+ die_errno("feed_pipe");
1755
+ if (ret) {
1756
+ /* done feeding */
1757
+ close(pp->children[i].process.in);
1758
+ pp->children[i].process.in = 0;
1759
+ }
1760
+ } else {
1761
+ /*
1762
+ * No feed_pipe means there is nothing to do, so
1763
+ * close the fd. Child input can be fed by other
1764
+ * methods, such as opts->path_to_stdin which
1765
+ * slurps a file via dup2, so clean up here.
1766
+ */
1767
+ close(pp->children[i].process.in);
1768
+ pp->children[i].process.in = 0;
1769
+ }
1770
+ }
1771
+
1772
+ /* Handle output reading (stderr) */
1773
if (child_is_working(&pp->children[i]) &&
1724
- pp->pfd[i].revents & (POLLIN | POLLHUP)) {
1774
+ pp->pfd[2*i].revents & (POLLIN | POLLHUP)) {
1775
int n = strbuf_read_once(&pp->children[i].err,
1776
pp->children[i].process.err, 0);
1777
if (n == 0) {
1864
1865
static void pp_handle_child_IO(struct parallel_processes *pp,
1866
const struct run_process_parallel_opts *opts,
1817
- int output_timeout)
1867
+ int timeout)
1868
{
1819
- /*
1820
- * First push input, if any (it might no-op), to child tasks to avoid them blocking
1821
- * after input. This also prevents deadlocks when ungrouping below, if a child blocks
1822
- * while the parent also waits for them to finish.
1823
- */
1824
- pp_buffer_stdin(pp, opts);
1825
-
1869
if (opts->ungroup) {
1870
+ pp_buffer_stdin(pp, opts);
1871
for (size_t i = 0; i < opts->processes; i++)
1872
if (child_is_ready_for_cleanup(&pp->children[i]))
1873
pp->children[i].state = GIT_CP_WAIT_CLEANUP;
1874
} else {
1831
- pp_buffer_stderr(pp, opts, output_timeout);
1875
+ pp_buffer_io(pp, opts, timeout);
1876
pp_output(pp);
1877
}
1878
}
1880
void run_processes_parallel(const struct run_process_parallel_opts *opts)
1881
{
1882
int i, code;
1839
- int output_timeout = 100;
1883
+ int timeout = 100;
1884
int spawn_cap = 4;
1885
struct parallel_processes_for_signal pp_sig;
1886
struct parallel_processes pp = {
1920
}
1921
if (!pp.nr_processes)
1922
break;
1879
- pp_handle_child_IO(&pp, opts, output_timeout);
1923
+ pp_handle_child_IO(&pp, opts, timeout);
1924
code = pp_collect_finished(&pp, opts);
1925
if (code) {
1926
pp.shutdown = 1;