{"thread":{"id":"24589","subject":"[remote-fd RFC PATCH] Rewrite bidirectional traffic loop","startedAt":"2010-07-31T04:38:10Z","lastAt":"2010-07-31T04:38:10Z","messageCount":1,"participants":["Ilari Liusvaara"],"isPatch":true,"patchVersion":1,"patchTotal":null},"messages":[{"id":"146835","messageId":"1280551090-19189-1-git-send-email-ilari.liusvaara@elisanet.fi","threadId":"24589","inReplyTo":null,"subject":"[remote-fd RFC PATCH] Rewrite bidirectional traffic loop","fromName":"Ilari Liusvaara","fromEmail":"ilari.liusvaara@elisanet.fi","sentAt":"2010-07-31T04:38:10Z","receivedAt":"2010-07-31T04:38:10Z","isPatch":true,"sender":{"key":"ilari.liusvaara@elisanet.fi","avatar":null},"body":"Rewrite bidirectional traffic loop to be clearer, fix some logic\nerrors that could lead into push failing or unneeded CPU usage,\nand support debugging mode (activated by setting $GIT_TRANSLOOP_DEBUG).\n\nSigned-off-by: Ilari Liusvaara <ilari.liusvaara@elisanet.fi>\n---\n transport-helper.c |  296 +++++++++++++++++++++++++++++++++++++---------------\n 1 files changed, 210 insertions(+), 86 deletions(-)\n\ndiff --git a/transport-helper.c b/transport-helper.c\nindex 3591e0d..0facf65 100644\n--- a/transport-helper.c\n+++ b/transport-helper.c\n@@ -865,130 +865,254 @@ int transport_helper_init(struct transport *transport, const char *name)\n \n \n #define BUFFERSIZE 4096\n+#define PBUFFERSIZE 8192\n \n-/* Copy data from stdin to output and from input to stdout. */\n-int bidirectional_transfer_loop(int input, int output)\n+/* Print bidirectional transfer loop debug message. */\n+static void transfer_debug(const char *fmt, ...)\n {\n-\tstruct pollfd polls[4];\n-\tchar in_buffer[BUFFERSIZE];\n-\tchar out_buffer[BUFFERSIZE];\n-\tsize_t in_buffer_use = 0;\n-\tsize_t out_buffer_use = 0;\n-\tint in_hup = 0;\n-\tint out_hup = 0;\n-\tint socket_mode = 0;\n-\tint input_index = 2;\n-\tint output_index = 3;\n-\tint poll_count = 4;\n+\tva_list args;\n+\tchar msgbuf[PBUFFERSIZE];\n+\tstatic int debug_enabled = -1;\n+\n+\tif (debug_enabled < 0)\n+\t\tdebug_enabled = getenv(\"GIT_TRANSLOOP_DEBUG\") ? 1 : 0;\n+\tif (!debug_enabled)\n+\t\treturn;\n+\n+\tsprintf(msgbuf, \"Transfer loop debugging: \");\n+\tva_start(args, fmt);\n+\tvsprintf(msgbuf + strlen(msgbuf), fmt, args);\n+\tva_end(args);\n+\tfprintf(stderr, \"%s\\n\", msgbuf);\n+}\n \n-\tif(input == output) {\n-\t\toutput_index = input_index;\n-\t\tpoll_count = 3;\n-\t\tsocket_mode = 1;\n+/* Load the parameters into poll structure. Return number of entries loaded */\n+static int load_poll_params(struct pollfd *polls, size_t inbufuse,\n+\tsize_t outbufuse, int in_hup, int out_hup, int in_closed,\n+\tint out_closed, int socket_mode, int input_fd, int output_fd)\n+{\n+\tint stdin_index = -1;\n+\tint stdout_index = -1;\n+\tint input_index = -1;\n+\tint output_index = -1;\n+\tint nextindex = 0;\n+\tint i;\n+\n+\t/*\n+\t * Inputs can't be waited at all if buffer is full since we can't\n+\t * do read on 0 bytes as it could do strange things.\n+\t */\n+\tif (!in_hup && inbufuse < BUFFERSIZE) {\n+\t\tstdin_index = nextindex++;\n+\t\tpolls[stdin_index].fd = 0;\n+\t\ttransfer_debug(\"Adding stdin to fds to wait for\");\n+\t}\n+\tif (!out_hup && outbufuse < BUFFERSIZE) {\n+\t\tinput_index = nextindex++;\n+\t\tpolls[input_index].fd = input_fd;\n+\t\ttransfer_debug(\"Adding remote input to fds to wait for\");\n+\t}\n+\tif (!out_closed && outbufuse > 0) {\n+\t\tstdout_index = nextindex++;\n+\t\tpolls[stdout_index].fd = 1;\n+\t\ttransfer_debug(\"Adding stdout to fds to wait for\");\n+\t}\n+\tif (!in_closed && inbufuse > 0) {\n+\t\tif (socket_mode && input_index >= 0)\n+\t\t\toutput_index = input_index;\n+\t\telse {\n+\t\t\toutput_index = nextindex++;\n+\t\t\tpolls[output_index].fd = output_fd;\n+\t\t}\n+\t\ttransfer_debug(\"Adding remote output to fds to wait for\");\n \t}\n \n-\twhile(in_buffer_use || out_buffer_use || !in_hup || !out_hup) {\n-\t\tint r, i;\n-\t\t/* Set up the poll and do it. */\n-\t\tpolls[0].fd = 0;\n-\t\tpolls[1].fd = 1;\n-\t\tpolls[input_index].fd = input;\n-\t\tpolls[output_index].fd = output;\n-\t\tfor(i = 0; i < 4; i++)\n-\t\t\tpolls[i].events = polls[i].revents = 0;\n-\n-\t\tif(in_buffer_use > 0)\n-\t\t\tpolls[output_index].events |= POLLOUT;\n-\t\tif(in_buffer_use < BUFFERSIZE && !in_hup)\n-\t\t\tpolls[0].events |= POLLIN;\n-\t\tif(out_buffer_use > 0)\n-\t\t\tpolls[1].events |= POLLOUT;\n-\t\tif(out_buffer_use < BUFFERSIZE && !out_hup)\n-\t\t\tpolls[input_index].events |= POLLIN;\n-\t\tr = poll(polls, poll_count, -1);\n-\t\tif(r < 0) {\n-\t\t\tif(errno == EWOULDBLOCK || errno == EAGAIN ||\n-\t\t\t\terrno == EINTR)\n-\t\t\t\tcontinue;\n-\t\t\tperror(\"poll failed\");\n-\t\t\treturn 1;\n-\t\t} else if(r == 0)\n-\t\t\tcontinue;\n+\tfor (i = 0; i < nextindex; i++)\n+\t\tpolls[i].events = polls[i].revents = 0;\n \n-\t\t/* Something interesting has happened... */\n-\t\tif(polls[0].revents & (POLLIN | POLLHUP)) {\n-\t\t\t/* Stdin is readable. */\n-\t\t\tr = read(0, in_buffer + in_buffer_use, BUFFERSIZE -\n-\t\t\t\tin_buffer_use);\n-\t\t\tif(r < 0 && errno != EWOULDBLOCK && errno != EAGAIN &&\n+\tif (stdin_index >= 0) {\n+\t\tpolls[stdin_index].events |= POLLIN;\n+\t\ttransfer_debug(\"Waiting for stdin to become readable\");\n+\t}\n+\tif (input_index >= 0) {\n+\t\tpolls[input_index].events |= POLLIN;\n+\t\ttransfer_debug(\"Waiting for remote input to become readable\");\n+\t}\n+\tif (stdout_index >= 0) {\n+\t\tpolls[stdout_index].events |= POLLOUT;\n+\t\ttransfer_debug(\"Waiting for stdout to become writable\");\n+\t}\n+\tif (output_index >= 0) {\n+\t\tpolls[output_index].events |= POLLOUT;\n+\t\ttransfer_debug(\"Waiting for remote output to become writable\");\n+\t}\n+\n+\t/* Return number of indexes assigned. */\n+\treturn nextindex;\n+}\n+\n+static int transfer_handle_events(struct pollfd* polls, char *in_buffer,\n+\tchar *out_buffer, size_t *in_buffer_use, size_t *out_buffer_use,\n+\tint *in_hup, int *out_hup, int *in_closed, int *out_closed,\n+\tint socket_mode, int poll_count, int input, int output)\n+{\n+\tint i, r;\n+\tfor(i = 0; i < poll_count; i++) {\n+\t\t/* Handle stdin. */\n+\t\tif (polls[i].fd == 0 && polls[i].revents & (POLLIN | POLLHUP)) {\n+\t\t\ttransfer_debug(\"stdin is readable\");\n+\t\t\tr = read(0, in_buffer + *in_buffer_use, BUFFERSIZE -\n+\t\t\t\t*in_buffer_use);\n+\t\t\tif (r < 0 && errno != EWOULDBLOCK && errno != EAGAIN &&\n \t\t\t\terrno != EINTR) {\n \t\t\t\tperror(\"read(git) failed\");\n \t\t\t\treturn 1;\n-\t\t\t} else if(r == 0) {\n-\t\t\t\tin_hup = 1;\n-\t\t\t\tif(!in_buffer_use) {\n-\t\t\t\t\tif(socket_mode)\n+\t\t\t} else if (r == 0) {\n+\t\t\t\ttransfer_debug(\"stdin EOF\");\n+\t\t\t\t*in_hup = 1;\n+\t\t\t\tif (!*in_buffer_use) {\n+\t\t\t\t\tif (socket_mode)\n \t\t\t\t\t\tshutdown(output, SHUT_WR);\n \t\t\t\t\telse\n \t\t\t\t\t\tclose(output);\n-\t\t\t\t}\n-\t\t\t} else\n-\t\t\t\tin_buffer_use += r;\n+\t\t\t\t\t*in_closed = 1;\n+\t\t\t\t\ttransfer_debug(\"Closed remote output\");\n+\t\t\t\t} else\n+\t\t\t\t\ttransfer_debug(\"Delaying remote output close because input buffer has data\");\n+\t\t\t} else if (r > 0) {\n+\t\t\t\t*in_buffer_use += r;\n+\t\t\t\ttransfer_debug(\"Read %i bytes from stdin (buffer now at %i)\", r, (int)*in_buffer_use);\n+\t\t\t}\n \t\t}\n \n-\t\tif(polls[input_index].revents & (POLLIN | POLLHUP)) {\n-\t\t\t/* Connection is readable. */\n-\t\t\tr = read(input, out_buffer + out_buffer_use,\n-\t\t\t\tBUFFERSIZE - out_buffer_use);\n-\t\t\tif(r < 0 && errno != EWOULDBLOCK && errno != EAGAIN &&\n+\t\t/* Handle remote end input. */\n+\t\tif (polls[i].fd == input &&\n+\t\t\tpolls[i].revents & (POLLIN | POLLHUP)) {\n+\t\t\ttransfer_debug(\"remote input is readable\");\n+\t\t\tr = read(input, out_buffer + *out_buffer_use,\n+\t\t\t\tBUFFERSIZE - *out_buffer_use);\n+\t\t\tif (r < 0 && errno != EWOULDBLOCK && errno != EAGAIN &&\n \t\t\t\terrno != EINTR) {\n \t\t\t\tperror(\"read(connection) failed\");\n \t\t\t\treturn 1;\n-\t\t\t} else if(r == 0) {\n-\t\t\t\tout_hup = 1;\n-\t\t\t\tif(!out_buffer_use)\n+\t\t\t} else if (r == 0) {\n+\t\t\t\ttransfer_debug(\"remote input EOF\");\n+\t\t\t\t*out_hup = 1;\n+\t\t\t\tif (!*out_buffer_use) {\n \t\t\t\t\tclose(1);\n-\t\t\t} else\n-\t\t\t\tout_buffer_use += r;\n+\t\t\t\t\t*out_closed = 1;\n+\t\t\t\t\ttransfer_debug(\"Closed stdout\");\n+\t\t\t\t} else\n+\t\t\t\t\ttransfer_debug(\"Delaying stdout close because output buffer has data\");\n+\n+\t\t\t} else if (r > 0) {\n+\t\t\t\t*out_buffer_use += r;\n+\t\t\t\ttransfer_debug(\"Read %i bytes from remote input (buffer now at %i)\", r, (int)*out_buffer_use);\n+\t\t\t}\n \t\t}\n \n-\t\tif(polls[1].revents & POLLOUT) {\n-\t\t\t/* Stdout is writable. */\n-\t\t\tr = write(1, out_buffer, out_buffer_use);\n-\t\t\tif(r < 0 && errno != EWOULDBLOCK && errno != EAGAIN &&\n+\t\t/* Handle stdout. */\n+\t\tif (polls[i].fd == 1 && polls[i].revents & POLLNVAL) {\n+\t\t\terror(\"Write pipe to Git unexpectedly closed.\");\n+\t\t\treturn 1;\n+\t\t}\n+\t\tif (polls[i].fd == 1 && polls[i].revents & POLLOUT) {\n+\t\t\ttransfer_debug(\"stdout is writable\");\n+\t\t\tr = write(1, out_buffer, *out_buffer_use);\n+\t\t\tif (r < 0 && errno != EWOULDBLOCK && errno != EAGAIN &&\n \t\t\t\terrno != EINTR) {\n \t\t\t\tperror(\"write(git) failed\");\n \t\t\t\treturn 1;\n-\t\t\t} else {\n-\t\t\t\tout_buffer_use -= r;\n-\t\t\t\tif(out_buffer_use > 0)\n+\t\t\t} else if (r > 0){\n+\t\t\t\t*out_buffer_use -= r;\n+\t\t\t\ttransfer_debug(\"Wrote %i bytes to stdout (buffer now at %i)\", r, (int)*out_buffer_use);\n+\t\t\t\tif (*out_buffer_use > 0)\n \t\t\t\t\tmemmove(out_buffer, out_buffer + r,\n-\t\t\t\t\t\tout_buffer_use);\n-\t\t\t\tif(out_hup && !out_buffer_use)\n+\t\t\t\t\t\t*out_buffer_use);\n+\t\t\t\tif (*out_hup && !*out_buffer_use) {\n \t\t\t\t\tclose(1);\n+\t\t\t\t\t*out_closed = 1;\n+\t\t\t\t\ttransfer_debug(\"Closed stdout\");\n+\t\t\t\t}\n \t\t\t}\n \t\t}\n \n-\t\tif(polls[output_index].revents & POLLOUT) {\n-\t\t\t/* Connection is writable. */\n-\t\t\tr = write(output, in_buffer, in_buffer_use);\n-\t\t\tif(r < 0 && errno != EWOULDBLOCK && errno != EAGAIN &&\n+\t\t/* Handle remote end output. */\n+\t\tif (polls[i].fd == output && polls[i].revents & POLLNVAL) {\n+\t\t\terror(\"Write pipe to remote end unexpectedly closed.\");\n+\t\t\treturn 1;\n+\t\t}\n+\t\tif (polls[i].fd == output && polls[i].revents & POLLOUT) {\n+\t\t\ttransfer_debug(\"remote output is writable\");\n+\t\t\tr = write(output, in_buffer, *in_buffer_use);\n+\t\t\tif (r < 0 && errno != EWOULDBLOCK && errno != EAGAIN &&\n \t\t\t\terrno != EINTR) {\n \t\t\t\tperror(\"write(connection) failed\");\n \t\t\t\treturn 1;\n-\t\t\t} else {\n-\t\t\t\tin_buffer_use -= r;\n-\t\t\t\tif(in_buffer_use > 0)\n+\t\t\t} else if (r > 0) {\n+\t\t\t\t*in_buffer_use -= r;\n+\t\t\t\ttransfer_debug(\"Wrote %i bytes to remote output (buffer now at %i)\", r, (int)*in_buffer_use);\n+\t\t\t\tif (*in_buffer_use > 0)\n \t\t\t\t\tmemmove(in_buffer, in_buffer + r,\n-\t\t\t\t\t\tin_buffer_use);\n-\t\t\t\tif(in_hup && !in_buffer_use) {\n-\t\t\t\t\tif(socket_mode)\n+\t\t\t\t\t\t*in_buffer_use);\n+\t\t\t\tif (*in_hup && !*in_buffer_use) {\n+\t\t\t\t\tif (socket_mode)\n \t\t\t\t\t\tshutdown(output, SHUT_WR);\n \t\t\t\t\telse\n \t\t\t\t\t\tclose(output);\n+\t\t\t\t\t*in_closed = 1;\n+\t\t\t\t\ttransfer_debug(\"Closed remote output\");\n \t\t\t\t}\n \t\t\t}\n \t\t}\n \t}\n \treturn 0;\n }\n+\n+/* Copy data from stdin to output and from input to stdout. */\n+int bidirectional_transfer_loop(int input, int output)\n+{\n+\tstruct pollfd polls[4];\n+\tchar in_buffer[BUFFERSIZE];\n+\tchar out_buffer[BUFFERSIZE];\n+\tsize_t in_buffer_use = 0;\n+\tsize_t out_buffer_use = 0;\n+\tint in_hup = 0;\n+\tint out_hup = 0;\n+\tint in_closed = 0;\n+\tint out_closed = 0;\n+\tint socket_mode = 0;\n+\tint poll_count = 4;\n+\n+\tif (input == output)\n+\t\tsocket_mode = 1;\n+\n+\twhile (1) {\n+\t\tint r;\n+\t\tpoll_count = load_poll_params(polls, in_buffer_use,\n+\t\t\tout_buffer_use, in_hup, out_hup, in_closed, out_closed,\n+\t\t\tsocket_mode, input, output);\n+\t\tif (!poll_count) {\n+\t\t\ttransfer_debug(\"Transfer done\");\n+\t\t\tbreak;\n+\t\t}\n+\t\ttransfer_debug(\"Waiting for %i file descriptors\", poll_count);\n+\t\tr = poll(polls, poll_count, -1);\n+\t\tif (r < 0) {\n+\t\t\tif (errno == EWOULDBLOCK || errno == EAGAIN ||\n+\t\t\t\terrno == EINTR)\n+\t\t\t\tcontinue;\n+\t\t\tperror(\"poll failed\");\n+\t\t\treturn 1;\n+\t\t} else if (r == 0)\n+\t\t\tcontinue;\n+\n+\t\tr = transfer_handle_events(polls, in_buffer, out_buffer,\n+\t\t\t&in_buffer_use, &out_buffer_use, &in_hup, &out_hup,\n+\t\t\t&in_closed, &out_closed, socket_mode, poll_count,\n+\t\t\tinput, output);\n+\t\tif (r)\n+\t\t\treturn r;\n+\t}\n+\treturn 0;\n+}\n-- \n1.7.2.1.9.g1ccab\n"}]}