diff --git a/Documentation/config/checkout.adoc b/Documentation/config/checkout.adoc index e35d21296978fe..45951bf38a5e3c 100644 --- a/Documentation/config/checkout.adoc +++ b/Documentation/config/checkout.adoc @@ -30,6 +30,11 @@ commands or functionality in the future. all commands that perform checkout. E.g. checkout, clone, reset, sparse-checkout, etc. + +On Windows the number of workers is capped at 62, because the `poll()` +emulation cannot wait on more worker pipes than that. A higher configured +value, including the logical core count on a machine with many cores, is +silently reduced to the cap. ++ NOTE: Parallel checkout usually delivers better performance for repositories located on SSDs or over NFS. For repositories on spinning disks and/or machines with a small number of cores, the default sequential checkout often performs diff --git a/Makefile b/Makefile index 87505e5df83ebf..17634413fc7907 100644 --- a/Makefile +++ b/Makefile @@ -1540,6 +1540,7 @@ CLAR_TEST_SUITES += u-odb-inmemory CLAR_TEST_SUITES += u-oid-array CLAR_TEST_SUITES += u-oidmap CLAR_TEST_SUITES += u-oidtree +CLAR_TEST_SUITES += u-poll CLAR_TEST_SUITES += u-prio-queue CLAR_TEST_SUITES += u-reftable-basics CLAR_TEST_SUITES += u-reftable-block diff --git a/builtin/fetch.c b/builtin/fetch.c index c1d7c672f4e0d8..48ea21ae28a945 100644 --- a/builtin/fetch.c +++ b/builtin/fetch.c @@ -2313,6 +2313,7 @@ static int fetch_multiple(struct string_list *list, int max_children, .tr2_label = "parallel/fetch", .processes = max_children, + .no_stdin_pipe = 1, .get_next_task = &fetch_next_remote, .start_failure = &fetch_failed_to_start, diff --git a/builtin/submodule--helper.c b/builtin/submodule--helper.c index 1cc82a134db22e..71a6650eca47d2 100644 --- a/builtin/submodule--helper.c +++ b/builtin/submodule--helper.c @@ -2914,6 +2914,7 @@ static int update_submodules(struct update_data *update_data) .tr2_label = "parallel/update", .processes = update_data->max_jobs, + .no_stdin_pipe = 1, .get_next_task = update_clone_get_next_task, .start_failure = update_clone_start_failure, diff --git a/compat/poll/poll.c b/compat/poll/poll.c index ea362b4a8e2340..1205f09ce9fbc8 100644 --- a/compat/poll/poll.c +++ b/compat/poll/poll.c @@ -303,6 +303,43 @@ compute_revents (int fd, int sought, fd_set *rfds, fd_set *wfds, fd_set *efds) } #endif /* !MinGW */ +#ifdef WIN32_NATIVE +/* POLL_MAX_DESCRIPTORS descriptors, plus hEvent and the QS_ALLINPUT message + queue, must fit in one MsgWaitForMultipleObjects call, and the collected + handles plus the NULL sentinel must fit in handle_array. */ +#if POLL_MAX_DESCRIPTORS + 2 > MAXIMUM_WAIT_OBJECTS +#error POLL_MAX_DESCRIPTORS exceeds MAXIMUM_WAIT_OBJECTS +#endif +#if POLL_MAX_DESCRIPTORS + 2 > FD_SETSIZE + 2 +#error POLL_MAX_DESCRIPTORS does not fit in handle_array +#endif + +/* Undo the WSAEventSelect() calls made for the first NFD descriptors. */ +static void +reset_socket_events (struct pollfd *pfd, nfds_t nfd) +{ + nfds_t i; + + for (i = 0; i < nfd; i++) + { + HANDLE h; + + if (pfd[i].fd < 0) + continue; + if (!(pfd[i].events & (POLLIN | POLLRDNORM | POLLOUT | POLLWRNORM | + POLLWRBAND | POLLPRI | POLLRDBAND))) + continue; + + h = (HANDLE) _get_osfhandle (pfd[i].fd); + if (h == NULL || h == INVALID_HANDLE_VALUE) + continue; + + if (IsSocketHandle (h)) + WSAEventSelect ((SOCKET) h, NULL, 0); + } +} +#endif + int poll (struct pollfd *pfd, nfds_t nfd, int timeout) { @@ -504,7 +541,16 @@ poll (struct pollfd *pfd, nfds_t nfd, int timeout) bits for the "wrong" direction. */ pfd[i].revents = win32_compute_revents (h, &sought); if (sought) - handle_array[nhandles++] = h; + { + /* hEvent occupies handle_array[0]. See POLL_MAX_DESCRIPTORS. */ + if (nhandles > POLL_MAX_DESCRIPTORS) + { + reset_socket_events (pfd, i); + errno = EINVAL; + return -1; + } + handle_array[nhandles++] = h; + } if (pfd[i].revents) timeout = 0; } diff --git a/compat/poll/poll.h b/compat/poll/poll.h index 1e1597360f4485..5c1169acc2c003 100644 --- a/compat/poll/poll.h +++ b/compat/poll/poll.h @@ -59,6 +59,22 @@ typedef unsigned long nfds_t; extern int poll (struct pollfd *pfd, nfds_t nfd, int timeout); +#if (defined _WIN32 || defined __WIN32__) && ! defined __CYGWIN__ +/* + * This poll() is emulated with MsgWaitForMultipleObjects(), which waits on at + * most MAXIMUM_WAIT_OBJECTS (64) objects. Two of those are never available for + * polled descriptors: poll() waits on its own event object, and QS_ALLINPUT + * adds the thread message queue. Sockets do not count, because they are all + * multiplexed onto that one event object; every other descriptor takes a wait + * slot of its own. + * + * Callers that poll one or more descriptors per child must keep the number of + * simultaneously live descriptors within this limit. Exceeding it fails with + * EINVAL. + */ +#define POLL_MAX_DESCRIPTORS 62 +#endif + /* Define INFTIM only if doing so conforms to POSIX. */ #if !defined (_POSIX_C_SOURCE) && !defined (_XOPEN_SOURCE) #define INFTIM (-1) diff --git a/compat/posix.h b/compat/posix.h index e2e794cad7d419..1a77b198aa5bc7 100644 --- a/compat/posix.h +++ b/compat/posix.h @@ -133,6 +133,16 @@ /* Pull the compat stuff */ #include #endif + +/* + * compat/poll defines POLL_MAX_DESCRIPTORS to the largest number of + * descriptors its poll() emulation can wait on. A native poll() has no such + * limit, so callers that fan out one descriptor per child can clamp against + * this unconditionally. + */ +#ifndef POLL_MAX_DESCRIPTORS +#define POLL_MAX_DESCRIPTORS INT_MAX +#endif #ifdef HAVE_BSD_SYSCTL #include #endif diff --git a/hook.c b/hook.c index d10eef4763c679..5bd0935bae307a 100644 --- a/hook.c +++ b/hook.c @@ -798,6 +798,7 @@ int run_hooks_opt(struct repository *r, const char *hook_name, .processes = jobs, .ungroup = jobs == 1, + .no_stdin_pipe = !options->feed_pipe, .get_next_task = pick_next_hook, .start_failure = notify_start_failure, diff --git a/parallel-checkout.c b/parallel-checkout.c index 1eb277a0fc0a55..4595cf4e8d9250 100644 --- a/parallel-checkout.c +++ b/parallel-checkout.c @@ -671,6 +671,13 @@ int run_parallel_checkout(struct checkout *state, int num_workers, int threshold if (parallel_checkout.nr < num_workers) num_workers = parallel_checkout.nr; + /* + * gather_results_from_workers() polls one pipe per worker, so the + * worker count must stay within what poll() can wait on. + */ + if (num_workers > POLL_MAX_DESCRIPTORS) + num_workers = POLL_MAX_DESCRIPTORS; + if (num_workers <= 1 || parallel_checkout.nr < threshold) { write_items_sequentially(state); } else { diff --git a/run-command.c b/run-command.c index e70a8a387b9042..99826b5442a1d9 100644 --- a/run-command.c +++ b/run-command.c @@ -1658,6 +1658,9 @@ static int pp_start_one(struct parallel_processes *pp, } return 1; } + if (opts->no_stdin_pipe && pp->children[i].process.in < 0) + BUG("get_next_task requested a stdin pipe despite " + "no_stdin_pipe"); if (!opts->ungroup) { pp->children[i].process.err = -1; pp->children[i].process.stdout_to_stderr = 1; @@ -1893,6 +1896,7 @@ void run_processes_parallel(const struct run_process_parallel_opts *opts) int i, code; int timeout = 100; int spawn_cap = 4; + size_t max_live; struct parallel_processes_for_signal pp_sig; struct parallel_processes pp = { .buffered_output = STRBUF_INIT, @@ -1902,6 +1906,20 @@ void run_processes_parallel(const struct run_process_parallel_opts *opts) const char *tr2_label = opts->tr2_label; const int do_trace2 = tr2_category && tr2_label; + /* + * Unless the caller handles its own output, pp_buffer_io() polls one + * output pipe per child and, unless excluded by no_stdin_pipe, may also + * poll an input pipe. Limit the number of live children so that all of + * their descriptors fit in one poll() call. + */ + max_live = opts->processes; + if (!opts->ungroup) { + size_t fds_per_process = opts->no_stdin_pipe ? 1 : 2; + + if (max_live > POLL_MAX_DESCRIPTORS / fds_per_process) + max_live = POLL_MAX_DESCRIPTORS / fds_per_process; + } + if (do_trace2) trace2_region_enter_printf(tr2_category, tr2_label, NULL, "max:%"PRIuMAX, @@ -1923,7 +1941,7 @@ void run_processes_parallel(const struct run_process_parallel_opts *opts) while (1) { for (i = 0; i < spawn_cap && !pp.shutdown && - pp.nr_processes < opts->processes; + pp.nr_processes < max_live; i++) { code = pp_start_one(&pp, opts); if (!code) diff --git a/run-command.h b/run-command.h index c2fad4f0f8380c..331c9391e8e2d5 100644 --- a/run-command.h +++ b/run-command.h @@ -488,6 +488,12 @@ struct run_process_parallel_opts */ unsigned int ungroup:1; + /** + * no_stdin_pipe: set if get_next_task will never request a pipe by + * setting child_process.in to -1. + */ + unsigned int no_stdin_pipe:1; + /** * get_next_task: See get_next_task_fn() above. This must be * specified. diff --git a/submodule.c b/submodule.c index fd91201a92d7b0..f40a82ee6aa33b 100644 --- a/submodule.c +++ b/submodule.c @@ -1835,6 +1835,7 @@ int fetch_submodules(struct repository *r, .tr2_label = "parallel/fetch", .processes = max_parallel_jobs, + .no_stdin_pipe = 1, .get_next_task = get_next_submodule, .start_failure = fetch_start_failure, diff --git a/t/helper/test-run-command.c b/t/helper/test-run-command.c index 4a56456894ccff..f943c8214d4dec 100644 --- a/t/helper/test-run-command.c +++ b/t/helper/test-run-command.c @@ -20,13 +20,14 @@ #include "wildmatch.h" static int number_callbacks; +static int max_callbacks = 4; static int parallel_next(struct child_process *cp, struct strbuf *err, void *cb, void **task_cb) { struct child_process *d = cb; - if (number_callbacks >= 4) + if (number_callbacks >= max_callbacks) return 0; strvec_pushv(&cp->args, d->args.v); @@ -195,6 +196,7 @@ static int testsuite(int argc, const char **argv) OPT_END() }; struct run_process_parallel_opts opts = { + .no_stdin_pipe = 1, .get_next_task = next_test, .start_failure = test_failed, .feed_pipe = test_stdin_pipe_feed, @@ -442,6 +444,16 @@ static int inherit_handle_child(void) int cmd__run_command(int argc, const char **argv) { struct child_process proc = CHILD_PROCESS_INIT; + const char * const parallel_usage[] = { + "test-tool run-command [] " + " [...]", + NULL + }; + struct option parallel_options[] = { + OPT_INTEGER_F(0, "tasks", &max_callbacks, + "number of tasks to generate", PARSE_OPT_NONEG), + OPT_END() + }; int jobs; int ret; struct run_process_parallel_opts opts = { @@ -495,17 +507,28 @@ int cmd__run_command(int argc, const char **argv) opts.ungroup = 1; } + argc = parse_options(argc - 1, argv + 1, NULL, parallel_options, + parallel_usage, PARSE_OPT_STOP_AT_NON_OPTION | + PARSE_OPT_KEEP_ARGV0); + if (argc < 3) + usage_with_options(parallel_usage, parallel_options); + if (max_callbacks < 0) + die("--tasks cannot be negative"); + jobs = atoi(argv[2]); strvec_clear(&proc.args); strvec_pushv(&proc.args, (const char **)argv + 3); if (!strcmp(argv[1], "run-command-parallel")) { + opts.no_stdin_pipe = 1; opts.get_next_task = parallel_next; opts.task_finished = task_finished_quiet; } else if (!strcmp(argv[1], "run-command-abort")) { + opts.no_stdin_pipe = 1; opts.get_next_task = parallel_next; opts.task_finished = task_finished; } else if (!strcmp(argv[1], "run-command-no-jobs")) { + opts.no_stdin_pipe = 1; opts.get_next_task = no_job; opts.task_finished = task_finished; } else if (!strcmp(argv[1], "run-command-stdin")) { diff --git a/t/meson.build b/t/meson.build index c86bf0de910cda..67bdabdfab4836 100644 --- a/t/meson.build +++ b/t/meson.build @@ -11,6 +11,7 @@ clar_test_suites = [ 'unit-tests/u-oid-array.c', 'unit-tests/u-oidmap.c', 'unit-tests/u-oidtree.c', + 'unit-tests/u-poll.c', 'unit-tests/u-prio-queue.c', 'unit-tests/u-reftable-basics.c', 'unit-tests/u-reftable-block.c', diff --git a/t/t0061-run-command.sh b/t/t0061-run-command.sh index 905e90e1f72541..a9958b66dc9fa9 100755 --- a/t/t0061-run-command.sh +++ b/t/t0061-run-command.sh @@ -164,6 +164,77 @@ test_expect_success 'run_command runs ungrouped in parallel with more tasks than test_line_count = 4 err ' +wait_for_line_count () { + expected=$1 && + file=$2 && + + for i in $(test_seq 1 100) + do + if test "$(wc -l <"$file")" -eq "$expected" + then + return 0 + fi && + sleep 0.1 + done && + return 1 +} + +cleanup_parallel () { + touch release + if test -n "$parallel_pid" + then + wait "$parallel_pid" + fi +} + +test_expect_success MINGW 'setup poll descriptor limit test' ' + write_script wait-for-release <<-\EOF + echo started >>"$1" + if test "$3" = stdin + then + while read line + do + : + done + fi + while ! test -e "$2" + do + sleep 0.1 + done + EOF +' + +test_expect_success MINGW 'run_command uses full poll limit without stdin' ' + : >started && + rm -f release && + test-tool run-command run-command-parallel --tasks=40 40 \ + ./wait-for-release "$PWD/started" "$PWD/release" \ + >out 2>err & + parallel_pid=$! && + test_when_finished cleanup_parallel && + wait_for_line_count 40 started && + touch release && + wait "$parallel_pid" && + parallel_pid= +' + +test_expect_success MINGW 'run_command limits children with stdin pipes' ' + : >started && + rm -f release && + test-tool run-command run-command-stdin --tasks=40 40 \ + ./wait-for-release "$PWD/started" "$PWD/release" stdin \ + >out 2>err & + parallel_pid=$! && + test_when_finished cleanup_parallel && + wait_for_line_count 31 started && + sleep 1 && + test_line_count = 31 started && + touch release && + wait "$parallel_pid" && + parallel_pid= && + test_line_count = 40 started +' + test_expect_success 'run_command listens to stdin' ' cat >expect <<-\EOF && preloaded output of a child diff --git a/t/t2080-parallel-checkout-basics.sh b/t/t2080-parallel-checkout-basics.sh index 7ad96cd5cd24a3..94d1f4bf1e7718 100755 --- a/t/t2080-parallel-checkout-basics.sh +++ b/t/t2080-parallel-checkout-basics.sh @@ -319,5 +319,41 @@ test_expect_success MINGW 'parallel checkout with fscache does not fail on new d test_cmp expect2 sub/deep/dir/file2.txt ) ' +# Windows has no native poll(). compat/poll emulates it with +# MsgWaitForMultipleObjects(), which cannot wait on more than +# MAXIMUM_WAIT_OBJECTS objects, so run_parallel_checkout() caps the worker +# count at MAXIMUM_WAIT_OBJECTS - 2. Without that cap, compat/poll collected +# one wait handle per polled worker pipe in a fixed-size stack array and +# smashed the stack. +# +# MAXIMUM_WAIT_OBJECTS is 64, hence the expected 62 below. The test is +# MINGW-only because the cap only exists there; on other platforms the +# requested 200 workers are used as-is. +test_expect_success MINGW 'checkout caps workers at the poll limit' ' + test_when_finished "rm -rf many-workers" && + git init many-workers && + ( + cd many-workers && + mkdir dir && + for i in $(test_seq 1 200) + do + echo "content $i" >dir/file$i || return 1 + done && + git add -A && + git commit -q -m base && + + git checkout -q -b other && + for i in $(test_seq 1 200) + do + echo "changed $i" >dir/file$i || return 1 + done && + git commit -q -a -m changed && + git checkout -q - + ) && + + set_checkout_config 200 1 && + test_checkout_workers 62 git -C many-workers checkout other && + verify_checkout many-workers +' test_done diff --git a/t/unit-tests/u-poll.c b/t/unit-tests/u-poll.c new file mode 100644 index 00000000000000..bc72135da28451 --- /dev/null +++ b/t/unit-tests/u-poll.c @@ -0,0 +1,163 @@ +#include "unit-test.h" + +#ifdef GIT_WINDOWS_NATIVE +static struct { + int pipes[POLL_MAX_DESCRIPTORS + 1][2]; + size_t nr_pipes; + int listener; + int sockets[2]; + WSAEVENT event; + int event_selected; +} poll_test; + +static void open_pipes(struct pollfd *fds, size_t nr) +{ + size_t i; + + cl_assert(nr <= ARRAY_SIZE(poll_test.pipes)); + for (i = 0; i < nr; i++) { + int *pipefd = poll_test.pipes[poll_test.nr_pipes]; + + cl_assert_equal_i(pipe(pipefd), 0); + poll_test.nr_pipes++; + fds[i].fd = pipefd[0]; + fds[i].events = POLLIN; + } +} + +static void create_socket_pair(void) +{ + struct sockaddr_in address = { + .sin_family = AF_INET, + .sin_addr.s_addr = htonl(INADDR_LOOPBACK), + }; + socklen_t address_length = sizeof(address); + + poll_test.listener = socket(AF_INET, SOCK_STREAM, 0); + cl_assert(poll_test.listener >= 0); + cl_assert_equal_i(bind(poll_test.listener, + (struct sockaddr *)&address, + sizeof(address)), 0); + cl_assert_equal_i(getsockname( + (SOCKET)_get_osfhandle(poll_test.listener), + (struct sockaddr *)&address, + &address_length), 0); + cl_assert_equal_i(listen(poll_test.listener, 1), 0); + + poll_test.sockets[0] = socket(AF_INET, SOCK_STREAM, 0); + cl_assert(poll_test.sockets[0] >= 0); + cl_assert_equal_i(connect(poll_test.sockets[0], + (struct sockaddr *)&address, + sizeof(address)), 0); + + poll_test.sockets[1] = accept(poll_test.listener, NULL, NULL); + cl_assert(poll_test.sockets[1] >= 0); + close(poll_test.listener); + poll_test.listener = -1; +} +#endif + +void test_poll__initialize(void) +{ +#ifdef GIT_WINDOWS_NATIVE + memset(&poll_test, 0, sizeof(poll_test)); + poll_test.listener = -1; + poll_test.sockets[0] = -1; + poll_test.sockets[1] = -1; + poll_test.event = WSA_INVALID_EVENT; +#endif +} + +void test_poll__cleanup(void) +{ +#ifdef GIT_WINDOWS_NATIVE + size_t i; + + if (poll_test.event_selected && poll_test.sockets[0] >= 0) + WSAEventSelect( + (SOCKET)_get_osfhandle(poll_test.sockets[0]), NULL, 0); + if (poll_test.event != WSA_INVALID_EVENT) + WSACloseEvent(poll_test.event); + if (poll_test.listener >= 0) + close(poll_test.listener); + for (i = 0; i < ARRAY_SIZE(poll_test.sockets); i++) + if (poll_test.sockets[i] >= 0) + close(poll_test.sockets[i]); + for (i = 0; i < poll_test.nr_pipes; i++) { + close(poll_test.pipes[i][0]); + close(poll_test.pipes[i][1]); + } +#endif +} + +void test_poll__limit(void) +{ +#ifdef GIT_WINDOWS_NATIVE + struct pollfd fds[POLL_MAX_DESCRIPTORS + 1] = { 0 }; + + open_pipes(fds, ARRAY_SIZE(fds)); + cl_assert_equal_i(poll(fds, POLL_MAX_DESCRIPTORS, 0), 0); + + errno = 0; + cl_assert_equal_i(poll(fds, ARRAY_SIZE(fds), 0), -1); + cl_assert_equal_i(errno, EINVAL); +#else + cl_skip(); +#endif +} + +void test_poll__sparse(void) +{ +#ifdef GIT_WINDOWS_NATIVE + struct pollfd fds[2 * POLL_MAX_DESCRIPTORS + 1] = { 0 }; + size_t i; + + open_pipes(fds, POLL_MAX_DESCRIPTORS); + create_socket_pair(); + + for (i = POLL_MAX_DESCRIPTORS; i > 0; i--) { + fds[2 * i - 1] = fds[i - 1]; + fds[2 * i - 2].fd = -1; + } + fds[2 * POLL_MAX_DESCRIPTORS].fd = poll_test.sockets[0]; + fds[2 * POLL_MAX_DESCRIPTORS].events = POLLIN; + + cl_assert_equal_i(poll(fds, ARRAY_SIZE(fds), 0), 0); +#else + cl_skip(); +#endif +} + +void test_poll__socket_cleanup(void) +{ +#ifdef GIT_WINDOWS_NATIVE + struct pollfd fds[POLL_MAX_DESCRIPTORS + 2] = { 0 }; + WSANETWORKEVENTS events; + SOCKET socket_handle; + + create_socket_pair(); + socket_handle = (SOCKET)_get_osfhandle(poll_test.sockets[0]); + poll_test.event = WSACreateEvent(); + cl_assert(poll_test.event != WSA_INVALID_EVENT); + cl_assert_equal_i(WSAEventSelect(socket_handle, poll_test.event, + FD_READ), 0); + poll_test.event_selected = 1; + + fds[0].fd = poll_test.sockets[0]; + open_pipes(fds + 1, POLL_MAX_DESCRIPTORS + 1); + + errno = 0; + cl_assert_equal_i(poll(fds, ARRAY_SIZE(fds), 0), -1); + cl_assert_equal_i(errno, EINVAL); + + cl_assert_equal_i(send( + (SOCKET)_get_osfhandle(poll_test.sockets[1]), + "x", 1, 0), 1); + cl_assert_equal_i(WaitForSingleObject(poll_test.event, 1000), + WAIT_OBJECT_0); + cl_assert_equal_i(WSAEnumNetworkEvents(socket_handle, poll_test.event, + &events), 0); +#else + cl_skip(); +#endif +}