From 68330b91840a440d25d477720f76fdd0badc25d0 Mon Sep 17 00:00:00 2001 From: Yosuke Shimizu Date: Fri, 18 Sep 2026 11:39:09 +0900 Subject: [PATCH] examples/echoserver, scripts: hold what a send did not take - WS_AppCtx stages each direction on its own: buffer, bufferIdx and bufferOff from appFd to the channel; fdBuffer, fdIdx and fdOff from the channel to appFd. thread_ctx_t drops channelBuffer and eofBuffer. - app_drain_to_channel() advances bufferOff by what wolfSSH_ChannelIdSend() took, and holds the rest on WS_WANT_WRITE, WS_WINDOW_FULL, WS_REKEYING, WS_CHANNEL_NOT_CONF and WS_CHAN_RXD, putting back the ssh->error from before that send for all but WS_WANT_WRITE. A zero or over-reported send ends the session. - app_pump_to_fd() reads the channel into fdBuffer and writes it to the pty or forward socket, both now nonblocking, keeping what would block. - ssh_worker() leaves appFd out of the read set while buffer holds bytes, puts it in the write set while fdBuffer does, and runs both pumps on every pass. - A pty read of zero ends the loop whatever errno holds. - app_echo_pump() echoes the shell channel through shellCtx.buffer and then hands each new chunk to process_bytes(). It replaces the EOF drain along with eofOff and eofRead. A negative read ends the session. - app_flush_fd() gives a connected forward up to APP_FLUSH_SECS without progress after the loop to take what fdBuffer still holds. - The agent keeps blocking writes, through app_write_all(), which retries an interrupted send. - A forward clears both of its buffers as it connects. - fwd-bulk.test phase 3 logs in with a password and echoes 32 MB through ssh -L while the client keeps sending but stops reading for 10 s. Phase 4 ends a forward while bytes are held for its target and checks the next forward's target gets only its own. Issue: F-10544 --- examples/echoserver/echoserver.c | 500 +++++++++++++++++++++---------- scripts/fwd-bulk.test | 286 +++++++++++++++++- 2 files changed, 629 insertions(+), 157 deletions(-) diff --git a/examples/echoserver/echoserver.c b/examples/echoserver/echoserver.c index 624701159..bb0b0633b 100644 --- a/examples/echoserver/echoserver.c +++ b/examples/echoserver/echoserver.c @@ -113,6 +113,7 @@ #define SOCKET_ECONNRESET ECONNRESET #define SOCKET_ECONNABORTED ECONNABORTED #define SOCKET_EWOULDBLOCK EWOULDBLOCK + #define SOCKET_EAGAIN EAGAIN #define SOCKET_EINTR EINTR #else #include @@ -120,6 +121,7 @@ #define SOCKET_ECONNRESET WSAECONNRESET #define SOCKET_ECONNABORTED WSAECONNABORTED #define SOCKET_EWOULDBLOCK WSAEWOULDBLOCK + #define SOCKET_EAGAIN WSAEWOULDBLOCK #define SOCKET_EINTR WSAEINTR #endif @@ -185,6 +187,14 @@ typedef struct WS_AppCtx { word32 channelId; WS_AppState state; byte buffer[EXAMPLE_BUFFER_SZ]; + /* Bytes staged in buffer and how many of them the channel has taken, + * with 0 <= bufferOff <= bufferIdx. Both survive a worker pass. */ + word32 bufferIdx; + word32 bufferOff; + /* Channel bytes waiting for appFd, and how many appFd has taken */ + byte fdBuffer[EXAMPLE_BUFFER_SZ]; + word32 fdIdx; + word32 fdOff; } WS_AppCtx; @@ -235,10 +245,6 @@ typedef struct { #ifdef WOLFSSH_SCP int doScp; #endif - byte channelBuffer[EXAMPLE_BUFFER_SZ]; - /* The EOF drain holds an unsent tail across worker passes, - * so it cannot share channelBuffer with the read path. */ - byte eofBuffer[EXAMPLE_BUFFER_SZ]; char statsBuffer[EXAMPLE_BUFFER_SZ]; } thread_ctx_t; @@ -876,6 +882,7 @@ static int wsShellStartCb(WOLFSSH_CHANNEL* channel, void* ctx) ShellChildCleanup(threadCtx); return 1; } + tcp_set_nonblocking(&threadCtx->shellCtx.appFd); /* Installed only now: the refusals above reap their own child. */ signal(SIGCHLD, ChildSig); @@ -1113,6 +1120,247 @@ static void buf_dump(unsigned char *buf, int len) #endif +/* Bytes appCtx still owes the channel. */ +static word32 app_staged(const WS_AppCtx* appCtx) +{ + return appCtx->bufferIdx - appCtx->bufferOff; +} + + +/* Hand the staged bytes to the channel, advancing bufferOff by however many + * it took. Returns 0 while the send is owed or done, negative to end the + * session. */ +static int app_drain_to_channel(WOLFSSH* ssh, WS_AppCtx* appCtx, + word32 channelId, int* wantWrite) +{ + int savedError; + int ret = 0; + int cnt; + + while (app_staged(appCtx) > 0) { + savedError = ssh->error; + cnt = wolfSSH_ChannelIdSend(ssh, channelId, + appCtx->buffer + appCtx->bufferOff, app_staged(appCtx)); + if (cnt > 0) { + if ((word32)cnt > app_staged(appCtx)) { + ret = WS_FATAL_ERROR; + break; + } + appCtx->bufferOff += (word32)cnt; + continue; + } + + if (cnt == WS_WANT_WRITE) { + *wantWrite = 1; + } + else if (cnt == WS_WINDOW_FULL || cnt == WS_REKEYING + || cnt == WS_CHANNEL_NOT_CONF || cnt == WS_CHAN_RXD) { + /* Owed, not failed: put back the code from before the send. A + * WS_WINDOW_FULL left here fails DoKexInit() on the peer's next + * KEXINIT and ends the session. */ + ssh->error = savedError; + } + else { + ret = (cnt < 0) ? cnt : WS_FATAL_ERROR; + } + break; + } + + if (app_staged(appCtx) == 0) { + appCtx->bufferIdx = 0; + appCtx->bufferOff = 0; + } + + return ret; +} + + +/* Loop the shell channel's buffered data back to it, staging in + * shellCtx.buffer, which echo mode leaves free. Sets dry once the channel + * holds nothing more. Returns 0, or negative to end the session. */ +static int app_echo_pump(WOLFSSH* ssh, thread_ctx_t* threadCtx, int* wantWrite, + int* dry) +{ + WS_AppCtx* appCtx = &threadCtx->shellCtx; + int cnt; + int ret; + + *dry = 0; + + for (;;) { + word32 freshSz = 0; + + if (app_staged(appCtx) == 0) { + cnt = wolfSSH_ChannelIdRead(ssh, appCtx->channelId, + appCtx->buffer, (word32)sizeof appCtx->buffer); + /* A negative read cannot be retried, so end the session. */ + if (cnt < 0) { + return cnt; + } + if (cnt == 0) { + *dry = 1; + break; + } + #ifdef SHELL_DEBUG + buf_dump(appCtx->buffer, cnt); + #endif + appCtx->bufferIdx = (word32)cnt; + appCtx->bufferOff = 0; + freshSz = (word32)cnt; + } + + ret = app_drain_to_channel(ssh, appCtx, appCtx->channelId, wantWrite); + if (ret < 0) { + return ret; + } + /* After the drain, so the echo precedes what the chunk asks for */ + if (freshSz > 0 && process_bytes(threadCtx, appCtx->buffer, freshSz)) { + ChildRunning = 0; + break; + } + if (app_staged(appCtx) > 0) { + break; + } + } + + return 0; +} + + +#ifdef WOLFSSH_AGENT + +/* Write every byte to a blocking socket, retrying what an interrupted + * write left behind. Returns bufSz, or -1 */ +static int app_write_all(WS_SOCKET_T fd, const byte* buf, word32 bufSz) +{ + word32 off = 0; + int cnt; + + while (off < bufSz) { + cnt = (int)send(fd, (const char*)buf + off, (int)(bufSz - off), 0); + + if (cnt > 0) { + off += (word32)cnt; + } + else if (cnt < 0 && SOCKET_ERRNO == SOCKET_EINTR) { + continue; + } + else { + return -1; + } + } + + return (int)bufSz; +} + +#endif /* WOLFSSH_AGENT */ + + +#if defined(WOLFSSH_SHELL) || defined(WOLFSSH_FWD) + +/* Bytes read from the channel that appFd has not taken yet. */ +static word32 app_fd_staged(const WS_AppCtx* appCtx) +{ + return appCtx->fdIdx - appCtx->fdOff; +} + + +/* Move what the channel holds into appFd, which must be nonblocking, keeping + * whatever a write would not take. Returns 0, or negative to end the + * session. */ +static int app_pump_to_fd(WOLFSSH* ssh, WS_AppCtx* appCtx, int isSocket) +{ + int cnt; + int err; + +#ifndef WOLFSSH_SHELL + (void)isSocket; +#endif + + for (;;) { + if (app_fd_staged(appCtx) == 0) { + appCtx->fdIdx = 0; + appCtx->fdOff = 0; + cnt = wolfSSH_ChannelIdRead(ssh, appCtx->channelId, + appCtx->fdBuffer, (word32)sizeof appCtx->fdBuffer); + if (cnt <= 0) { + return cnt; + } + #ifdef SHELL_DEBUG + buf_dump(appCtx->fdBuffer, cnt); + #endif + appCtx->fdIdx = (word32)cnt; + } + +#ifdef WOLFSSH_SHELL + if (!isSocket) { + cnt = (int)write(appCtx->appFd, + appCtx->fdBuffer + appCtx->fdOff, + app_fd_staged(appCtx)); + } + else { + cnt = (int)send(appCtx->appFd, + (const char*)appCtx->fdBuffer + appCtx->fdOff, + (int)app_fd_staged(appCtx), 0); + } +#else + cnt = (int)send(appCtx->appFd, + (const char*)appCtx->fdBuffer + appCtx->fdOff, + (int)app_fd_staged(appCtx), 0); +#endif + + if (cnt > 0) { + appCtx->fdOff += (word32)cnt; + continue; + } + err = SOCKET_ERRNO; + if (cnt < 0 && err == SOCKET_EINTR) { + continue; + } + /* appFd is full; its entry in the write set wakes the retry. */ + if (cnt < 0 && (err == SOCKET_EWOULDBLOCK || err == SOCKET_EAGAIN)) { + return 0; + } + return WS_FATAL_ERROR; + } +} + +#endif /* WOLFSSH_SHELL || WOLFSSH_FWD */ + + +#ifdef WOLFSSH_FWD + +/* Seconds a stalled forward gets, once the session is over, to take what the + * channel still holds. */ +#define APP_FLUSH_SECS 5 + +/* A peer that ends the session right after its data would otherwise lose + * what the forward had not written yet. */ +static void app_flush_fd(WOLFSSH* ssh, WS_AppCtx* appCtx) +{ + int waited = 0; + int sel; + + for (;;) { + if (app_pump_to_fd(ssh, appCtx, 1) < 0 + || app_fd_staged(appCtx) == 0) { + break; + } + sel = tcp_select_write(appCtx->appFd, 1); + if (sel == WS_SELECT_TIMEOUT) { + if (++waited >= APP_FLUSH_SECS) { + break; + } + } + else if (sel != WS_SELECT_SEND_READY) { + break; + } + } +} + +#endif /* WOLFSSH_FWD */ + + static int ssh_worker(thread_ctx_t* threadCtx) { WOLFSSH* ssh; @@ -1122,9 +1370,6 @@ static int ssh_worker(thread_ctx_t* threadCtx) * still leaves through the cleanup below it. */ int workerRet = 0; int eofAnswered = 0; - /* Held across passes with 0 <= eofOff <= eofRead. */ - int eofRead = 0; - int eofOff = 0; /* Without a shell there is no child to outlive the peer's EOF, and the * read path echoes unconditionally. */ int echoOnly = 1; @@ -1188,7 +1433,6 @@ static int ssh_worker(thread_ctx_t* threadCtx) #endif #ifdef WOLFSSH_FWD WS_SOCKET_T fwdFd = -1; - word32 fwdBufferIdx = 0; #endif ChildRunning = 1; @@ -1197,9 +1441,9 @@ static int ssh_worker(thread_ctx_t* threadCtx) fd_set readFds; fd_set writeFds; int writable; + int writeArmed; WS_SOCKET_T maxFd; int cnt_r; - int cnt_w; FD_ZERO(&readFds); FD_SET(sshFd, &readFds); @@ -1227,13 +1471,22 @@ static int ssh_worker(thread_ctx_t* threadCtx) wantWrite = 1; FD_ZERO(&writeFds); + writeArmed = wantWrite; if (wantWrite) FD_SET(sshFd, &writeFds); + /* Keep appFd out of the read set while buffer holds bytes for the + * channel, and in the write set while fdBuffer holds bytes for + * appFd. */ #ifdef WOLFSSH_SHELL if (threadCtx->shellCtx.state == APP_STATE_CONNECTED && threadCtx->shellCtx.appFd >= 0) { - FD_SET(threadCtx->shellCtx.appFd, &readFds); + if (app_staged(&threadCtx->shellCtx) == 0) + FD_SET(threadCtx->shellCtx.appFd, &readFds); + if (app_fd_staged(&threadCtx->shellCtx) > 0) { + FD_SET(threadCtx->shellCtx.appFd, &writeFds); + writeArmed = 1; + } if (threadCtx->shellCtx.appFd > maxFd) maxFd = threadCtx->shellCtx.appFd; } @@ -1248,7 +1501,8 @@ static int ssh_worker(thread_ctx_t* threadCtx) maxFd = threadCtx->agentCtx.listenFd; } if (agentFd >= 0 - && threadCtx->agentCtx.state == APP_STATE_CONNECTED) { + && threadCtx->agentCtx.state == APP_STATE_CONNECTED + && app_staged(&threadCtx->agentCtx) == 0) { FD_SET(agentFd, &readFds); if (agentFd > maxFd) maxFd = agentFd; @@ -1265,14 +1519,19 @@ static int ssh_worker(thread_ctx_t* threadCtx) } if (fwdFd >= 0 && threadCtx->fwdCtx.state == APP_STATE_CONNECTED) { - FD_SET(fwdFd, &readFds); + if (app_staged(&threadCtx->fwdCtx) == 0) + FD_SET(fwdFd, &readFds); + if (app_fd_staged(&threadCtx->fwdCtx) > 0) { + FD_SET(fwdFd, &writeFds); + writeArmed = 1; + } if (fwdFd > maxFd) maxFd = fwdFd; } #endif /* WOLFSSH_FWD */ rc = select((int)maxFd + 1, &readFds, - wantWrite ? &writeFds : NULL, NULL, NULL); + writeArmed ? &writeFds : NULL, NULL, NULL); if (rc == -1) { break; } @@ -1325,39 +1584,11 @@ static int ssh_worker(thread_ctx_t* threadCtx) threadCtx->shellCtx.channelId, WS_CHANNEL_ID_SELF); if (eofChannel != NULL && wolfSSH_ChannelGetEof(eofChannel)) { - int eofSent; int eofDrained = 0; - for (;;) { - /* A send is bounded by the peer's window and - * packet size, so a short one is normal. Read - * the next chunk only once the last one is out: - * the read consumed it from the channel, so its - * tail cannot be dropped. */ - if (eofOff == eofRead) { - int eofRxd; - - eofOff = eofRead = 0; - eofRxd = wolfSSH_ChannelIdRead(ssh, - threadCtx->shellCtx.channelId, - threadCtx->eofBuffer, - sizeof threadCtx->eofBuffer); - /* A negative read is a rekey or a stalled - * channel, not a drained one. */ - if (eofRxd <= 0) { - eofDrained = (eofRxd == 0); - break; - } - eofRead = eofRxd; - } - - eofSent = wolfSSH_ChannelIdSend(ssh, - threadCtx->shellCtx.channelId, - threadCtx->eofBuffer + eofOff, - eofRead - eofOff); - if (eofSent <= 0) - break; - eofOff += eofSent; + if (app_echo_pump(ssh, threadCtx, &wantWrite, + &eofDrained) < 0) { + break; } /* Only an emptied channel earns the EOF; anything @@ -1385,57 +1616,13 @@ static int ssh_worker(thread_ctx_t* threadCtx) * wolfSSH_ChannelIdRead() has no isKeying gate; the window * credit it owes is parked until the rekey finishes. */ if (rc == WS_CHAN_RXD || rc == WS_REKEYING) { - if (threadCtx->shellCtx.state == APP_STATE_CONNECTED && - lastChannel == threadCtx->shellCtx.channelId) { - cnt_r = wolfSSH_ChannelIdRead(ssh, - threadCtx->shellCtx.channelId, - threadCtx->channelBuffer, - sizeof threadCtx->channelBuffer); - if (cnt_r <= 0) { - /* Nothing was buffered. Only an actual data - * report makes that a failure. */ - if (rc == WS_REKEYING && cnt_r == 0) - continue; - break; - } - #ifdef SHELL_DEBUG - buf_dump(threadCtx->channelBuffer, cnt_r); - #endif - #ifdef WOLFSSH_SHELL - if (!threadCtx->echo) { - cnt_w = (int)write( - threadCtx->shellCtx.appFd, - threadCtx->channelBuffer, cnt_r); - } - else { - cnt_w = wolfSSH_ChannelIdSend(ssh, - threadCtx->shellCtx.channelId, - threadCtx->channelBuffer, cnt_r); - if (cnt_r > 0) { - int doStop = process_bytes(threadCtx, - threadCtx->channelBuffer, - cnt_r); - ChildRunning = !doStop; - } - } - #else - cnt_w = wolfSSH_ChannelIdSend(ssh, - threadCtx->shellCtx.channelId, - threadCtx->channelBuffer, cnt_r); - if (cnt_r > 0) { - int doStop = process_bytes(threadCtx, - threadCtx->channelBuffer, cnt_r); - ChildRunning = !doStop; - } - #endif - if (cnt_w <= 0) - break; - } + /* The session and the forward are drained by their + * pumps below, on every pass. */ #ifdef WOLFSSH_AGENT if (lastChannel == agentChannelId) { cnt_r = wolfSSH_ChannelIdRead(ssh, agentChannelId, - threadCtx->channelBuffer, - sizeof threadCtx->channelBuffer); + threadCtx->agentCtx.fdBuffer, + sizeof threadCtx->agentCtx.fdBuffer); if (cnt_r <= 0) { /* Nothing was buffered. Only an actual data * report makes that a failure. */ @@ -1444,35 +1631,11 @@ static int ssh_worker(thread_ctx_t* threadCtx) break; } #ifdef SHELL_DEBUG - buf_dump(threadCtx->channelBuffer, cnt_r); + buf_dump(threadCtx->agentCtx.fdBuffer, cnt_r); #endif - cnt_w = (int)send(agentFd, - threadCtx->channelBuffer, cnt_r, 0); - if (cnt_w <= 0) - break; - } - #endif - #ifdef WOLFSSH_FWD - if (threadCtx->fwdCtx.state == APP_STATE_CONNECTED && - lastChannel == threadCtx->fwdCtx.channelId) { - - cnt_r = wolfSSH_ChannelIdRead(ssh, - threadCtx->fwdCtx.channelId, - threadCtx->channelBuffer, - sizeof threadCtx->channelBuffer); - if (cnt_r <= 0) { - /* Nothing was buffered. Only an actual data - * report makes that a failure. */ - if (rc == WS_REKEYING && cnt_r == 0) - continue; - break; - } - #ifdef SHELL_DEBUG - buf_dump(threadCtx->channelBuffer, cnt_r); - #endif - cnt_w = (int)send(fwdFd, threadCtx->channelBuffer, - cnt_r, 0); - if (cnt_w <= 0) + if (app_write_all(agentFd, + threadCtx->agentCtx.fdBuffer, + (word32)cnt_r) < 0) break; } #endif @@ -1544,7 +1707,7 @@ static int ssh_worker(thread_ctx_t* threadCtx) /* This read will return 0 on EOF */ if (cnt_r <= 0) { int err = errno; - if (err != EAGAIN) { + if (cnt_r == 0 || err != EAGAIN) { #ifdef SHELL_DEBUG printf("Break:read childFd returns %d: " "errno =%x\n", @@ -1557,17 +1720,36 @@ static int ssh_worker(thread_ctx_t* threadCtx) #ifdef SHELL_DEBUG buf_dump(threadCtx->shellCtx.buffer, cnt_r); #endif - if (cnt_r > 0) { - cnt_w = wolfSSH_ChannelIdSend(ssh, - threadCtx->shellCtx.channelId, - threadCtx->shellCtx.buffer, cnt_r); - if (cnt_w < 0) - break; - } + threadCtx->shellCtx.bufferIdx = (word32)cnt_r; + threadCtx->shellCtx.bufferOff = 0; } } } #endif /* WOLFSSH_SHELL */ + /* On every pass, whatever woke it: send buffer to the channel and + * the channel's data to appFd, or echo it in echo mode. */ + if (threadCtx->shellCtx.state == APP_STATE_CONNECTED) { + int echoDry; + + if (!echoOnly) { + if (app_drain_to_channel(ssh, &threadCtx->shellCtx, + threadCtx->shellCtx.channelId, + &wantWrite) < 0) { + break; + } + #ifdef WOLFSSH_SHELL + if (threadCtx->shellCtx.appFd >= 0 + && app_pump_to_fd(ssh, &threadCtx->shellCtx, + 0) < 0) { + break; + } + #endif + } + else if (app_echo_pump(ssh, threadCtx, &wantWrite, + &echoDry) < 0) { + break; + } + } #ifdef WOLFSSH_AGENT if (agentFd >= 0 && threadCtx->agentCtx.state == APP_STATE_CONNECTED) { @@ -1609,13 +1791,14 @@ static int ssh_worker(thread_ctx_t* threadCtx) #ifdef SHELL_DEBUG buf_dump(threadCtx->agentCtx.buffer, cnt_r); #endif - cnt_w = wolfSSH_ChannelIdSend(ssh, agentChannelId, - threadCtx->agentCtx.buffer, cnt_r); - if (cnt_w <= 0) { - break; - } + threadCtx->agentCtx.bufferIdx = (word32)cnt_r; + threadCtx->agentCtx.bufferOff = 0; } } + if (app_drain_to_channel(ssh, &threadCtx->agentCtx, + agentChannelId, &wantWrite) < 0) { + break; + } } if (threadCtx->agentCtx.state == APP_STATE_LISTEN && threadCtx->agentCtx.listenFd >= 0) { @@ -1631,6 +1814,8 @@ static int ssh_worker(thread_ctx_t* threadCtx) } } else { + threadCtx->agentCtx.bufferIdx = 0; + threadCtx->agentCtx.bufferOff = 0; threadCtx->agentCtx.state = APP_STATE_CONNECTED; threadCtx->agentCtx.appFd = agentFd; } @@ -1644,9 +1829,8 @@ static int ssh_worker(thread_ctx_t* threadCtx) #ifdef SHELL_DEBUG printf("fwdFd set in readfd\n"); #endif - cnt_r = (int)recv(fwdFd, - threadCtx->fwdCtx.buffer + fwdBufferIdx, - sizeof threadCtx->fwdCtx.buffer - fwdBufferIdx, 0); + cnt_r = (int)recv(fwdFd, threadCtx->fwdCtx.buffer, + sizeof threadCtx->fwdCtx.buffer, 0); if (cnt_r == 0) { /* Read zero-returned. Socket is closed. Go back to listening. */ @@ -1677,31 +1861,25 @@ static int ssh_worker(thread_ctx_t* threadCtx) threadCtx->fwdCtx.state = APP_STATE_LISTEN; continue; } - break; + if (err != SOCKET_EWOULDBLOCK + && err != SOCKET_EAGAIN) { + break; + } } else { #ifdef SHELL_DEBUG buf_dump(threadCtx->fwdCtx.buffer, cnt_r); #endif - fwdBufferIdx += cnt_r; + threadCtx->fwdCtx.bufferIdx = (word32)cnt_r; + threadCtx->fwdCtx.bufferOff = 0; } } - if (fwdBufferIdx > 0) { - cnt_w = wolfSSH_ChannelIdSend(ssh, - threadCtx->fwdCtx.channelId, - threadCtx->fwdCtx.buffer, fwdBufferIdx); - if (cnt_w > 0) { - fwdBufferIdx = 0; - } - else if (cnt_w == WS_CHANNEL_NOT_CONF || - cnt_w == WS_CHAN_RXD) { - #ifdef SHELL_DEBUG - printf("Waiting for channel open confirmation.\n"); - #endif - } - else { - break; - } + if (app_drain_to_channel(ssh, &threadCtx->fwdCtx, + threadCtx->fwdCtx.channelId, &wantWrite) < 0) { + break; + } + if (app_pump_to_fd(ssh, &threadCtx->fwdCtx, 1) < 0) { + break; } } if (threadCtx->fwdCtx.state == APP_STATE_LISTEN @@ -1723,6 +1901,7 @@ static int ssh_worker(thread_ctx_t* threadCtx) const char* out = NULL; char addr[200]; + tcp_set_nonblocking(&fwdFd); threadCtx->fwdCtx.state = APP_STATE_CONNECT; threadCtx->fwdCtx.appFd = fwdFd; originAddrSz = sizeof originAddr; @@ -1763,6 +1942,10 @@ static int ssh_worker(thread_ctx_t* threadCtx) threadCtx->fwdCbCtx.originName, threadCtx->fwdCbCtx.originPort); if (newChannel != NULL) { + threadCtx->fwdCtx.bufferIdx = 0; + threadCtx->fwdCtx.bufferOff = 0; + threadCtx->fwdCtx.fdIdx = 0; + threadCtx->fwdCtx.fdOff = 0; threadCtx->fwdCtx.state = APP_STATE_CONNECTED; } } @@ -1771,6 +1954,11 @@ static int ssh_worker(thread_ctx_t* threadCtx) threadCtx->fwdCbCtx.hostPort); if (fwdFd > 0) { + tcp_set_nonblocking(&fwdFd); + threadCtx->fwdCtx.bufferIdx = 0; + threadCtx->fwdCtx.bufferOff = 0; + threadCtx->fwdCtx.fdIdx = 0; + threadCtx->fwdCtx.fdOff = 0; threadCtx->fwdCtx.appFd = fwdFd; threadCtx->fwdCtx.state = APP_STATE_CONNECTED; threadCtx->fwdCbCtx.isDirect = 0; @@ -1778,6 +1966,10 @@ static int ssh_worker(thread_ctx_t* threadCtx) } #endif /* WOLFSSH_FWD */ } +#ifdef WOLFSSH_FWD + if (fwdFd >= 0 && threadCtx->fwdCtx.state == APP_STATE_CONNECTED) + app_flush_fd(ssh, &threadCtx->fwdCtx); +#endif #ifdef WOLFSSH_SHELL ShellChildCleanup(threadCtx); #endif diff --git a/scripts/fwd-bulk.test b/scripts/fwd-bulk.test index b327c38cd..50db3b03f 100755 --- a/scripts/fwd-bulk.test +++ b/scripts/fwd-bulk.test @@ -13,7 +13,10 @@ # windows long and compares the bytes that come out. Phase 2 ends a transfer # one window plus a short tail in, the point where the tail is still sitting # in portfwd's buffer waiting on window credit when the local socket reports -# end-of-input. +# end-of-input. Phase 3 echoes the payload back while the client keeps sending +# but stops reading for a while. Phase 4 ends a forward while the echoserver +# still holds bytes for its target, then checks that the next forward's target +# gets only its own bytes. # # The listening nc takes its stdin from a long sleep on purpose. Reading # end-of-input on stdin makes nc close the connection, which truncates the @@ -31,6 +34,13 @@ tail_nc_server_pid=$no_pid tail_server_pid=$no_pid tail_portfwd_pid=$no_pid tail_nc_client_pid=$no_pid +echo_target_pid=$no_pid +echo_server_pid=$no_pid +echo_ssh_hold_pid=$no_pid +echo_ssh_pid=$no_pid +echo_client_pid=$no_pid +reuse_target_pid=$no_pid +reuse_client_pid=$no_pid work_dir="`pwd`/wolfssh_fwd_bulk$$" ready_file="$work_dir/ready" fwd_ready_file="$work_dir/fwd_ready" @@ -43,6 +53,19 @@ tail_payload="$work_dir/tail_payload" tail_received="$work_dir/tail_received" tail_server_log="$work_dir/tail_server.log" tail_portfwd_log="$work_dir/tail_portfwd.log" +echo_payload="$work_dir/echo_payload" +echo_received="$work_dir/echo_received" +echo_target_ready="$work_dir/echo_target_ready" +echo_ready_file="$work_dir/echo_ready" +echo_ssh_fifo="$work_dir/echo_ssh_hold" +echo_askpass="$work_dir/echo_askpass" +echo_server_log="$work_dir/echo_server.log" +echo_ssh_log="$work_dir/echo_ssh.log" +reuse_target_ready="$work_dir/reuse_target_ready" +reuse_stalled="$work_dir/reuse_stalled" +reuse_closed="$work_dir/reuse_closed" +reuse_expected="$work_dir/reuse_expected" +reuse_received="$work_dir/reuse_received" # Several times the 128K default window, so the transfer cannot finish # without the window being credited back at least once. The size the # receiver is held to is read back from the file dd made. @@ -69,6 +92,21 @@ tail_got=0 transfer_limit=90 # Consecutive seconds with no new bytes before calling it stalled. stall_limit=15 +# Phase 3. Well past the point a deadlocked forward stops (under 9 MB on the +# hosts measured), with a stall long enough to fill the socket buffers both +# ways. The phase's own stall limit has to outlast that pause. +echo_payload_blocks=32000 +echo_payload_size=0 +echo_pause_after=1000000 +echo_pause_secs=10 +echo_stall_limit=20 +echo_entry_port=0 +echo_target_port=0 +# Phase 4. The second connection's payload, and seconds to wait for it. +reuse_message="second connection" +reuse_limit=30 +reuse_entry_port=0 +reuse_target_port=0 [ ! -x "`command -v nc`" ] && echo "nc doesn't exist, skipping" && exit 77 [ ! -x ./examples/echoserver/echoserver ] \ @@ -91,7 +129,10 @@ echo "$WOLFSSH_OPTIONS" | grep -qx "TEST_BLOCK" \ do_cleanup() { for pid in $nc_client_pid $portfwd_pid $server_pid $nc_server_pid \ $hold_pid $tail_nc_client_pid $tail_portfwd_pid \ - $tail_server_pid $tail_nc_server_pid $tail_hold_pid + $tail_server_pid $tail_nc_server_pid $tail_hold_pid \ + $echo_client_pid $echo_ssh_pid $echo_ssh_hold_pid \ + $echo_server_pid $echo_target_pid $reuse_client_pid \ + $reuse_target_pid do if [ "$pid" != "$no_pid" ] then @@ -105,7 +146,7 @@ do_cleanup() { # The failure here is usually just a byte count; the logs name the cause. do_dump_logs() { for log in "$server_log" "$portfwd_log" "$tail_server_log" \ - "$tail_portfwd_log" + "$tail_portfwd_log" "$echo_server_log" "$echo_ssh_log" do [ -f "$log" ] || continue echo "--- `basename "$log"` ---" @@ -326,5 +367,244 @@ do done echo "ended $tail_attempts transfers on a window boundary" + +# --- Phase 3: echo both ways while the client stops reading ----------------- +# Both clients have to keep sending while they stop reading. portfwd and nc +# stop both ways once their output blocks, so this uses OpenSSH and python. +if [ ! -x "`command -v python3`" ] || [ ! -x "`command -v ssh`" ] +then + echo "skipping phases 3 and 4: they need python3 and ssh" + do_cleanup + exit 0 +fi + +echo_entry_port=`expr 28000 + \( $$ % 1000 \)` +echo_target_port=`expr 29000 + \( $$ % 1000 \)` +reuse_entry_port=`expr 30000 + \( $$ % 1000 \)` +reuse_target_port=`expr 31000 + \( $$ % 1000 \)` + +dd if=/dev/urandom of="$echo_payload" bs=1000 count=$echo_payload_blocks \ + 2>/dev/null || do_fail "couldn't make the payload" +echo_payload_size=`wc -c < "$echo_payload" | tr -d ' '` + +# Reads the next chunk only once the last one has gone back, so it stops +# reading when it cannot write. Once the whole payload is back it ends its +# side, and exits when the echoserver has closed the forward. +python3 -c ' +import socket, sys +s = socket.socket() +s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) +s.bind(("127.0.0.1", int(sys.argv[1]))) +s.listen(1) +open(sys.argv[2], "w").write("up") +c = s.accept()[0] +left = int(sys.argv[3]) +while left > 0: + d = c.recv(min(65536, left)) + if not d: + break + c.sendall(d) + left -= len(d) +c.shutdown(socket.SHUT_WR) +while c.recv(65536): + pass +' $echo_target_port "$echo_target_ready" $echo_payload_size 2>/dev/null & +echo_target_pid=$! + +counter=0 +while [ ! -s "$echo_target_ready" ] && [ "$counter" -lt 30 ]; do + sleep 0.1 + counter=`expr $counter + 1` +done +[ -s "$echo_target_ready" ] \ + || do_fail "couldn't listen on port $echo_target_port, is it in use?" + +./examples/echoserver/echoserver -1 -f -R "$echo_ready_file" \ + > "$echo_server_log" 2>&1 & +echo_server_pid=$! + +counter=0 +while [ ! -s "$echo_ready_file" ] && [ "$counter" -lt 20 ]; do + sleep 0.1 + counter=`expr $counter + 1` +done +[ -s "$echo_ready_file" ] || do_fail "no ready file, echoserver didn't start" + +# The echoserver only serves a forward alongside a session. Hold the +# session's stdin open: its end-of-input would end the session and take the +# forward with it. Log in with a password: a user key needs an algorithm the +# build may leave out. +printf '#!/bin/sh\necho upthehill\n' > "$echo_askpass" \ + && chmod 700 "$echo_askpass" \ + || do_fail "couldn't stage the password helper" +mkfifo "$echo_ssh_fifo" || do_fail "couldn't make the fifo" +sleep 300 > "$echo_ssh_fifo" 2>/dev/null & +echo_ssh_hold_pid=$! +SSH_ASKPASS="$echo_askpass" SSH_ASKPASS_REQUIRE=force DISPLAY=:0 \ +ssh -v -T -F /dev/null -o PreferredAuthentications=password \ + -o NumberOfPasswordPrompts=1 \ + -o StrictHostKeyChecking=no -o UserKnownHostsFile=/dev/null \ + -o ExitOnForwardFailure=yes -p `cat "$echo_ready_file"` \ + -L $echo_entry_port:127.0.0.1:$echo_target_port \ + -L $reuse_entry_port:127.0.0.1:$reuse_target_port jill@127.0.0.1 \ + < "$echo_ssh_fifo" > /dev/null 2> "$echo_ssh_log" & +echo_ssh_pid=$! + +# Connecting to check would take the target's one connection, so wait for +# ssh to say the listener is up instead. +counter=0 +while ! grep -q "Local forwarding listening on" "$echo_ssh_log" 2>/dev/null \ + && [ "$counter" -lt 50 ]; do + sleep 0.1 + counter=`expr $counter + 1` +done +grep -q "Local forwarding listening on" "$echo_ssh_log" \ + || do_fail "ssh didn't set up the forward" + +python3 -c ' +import socket, sys, threading, time +s = socket.create_connection(("127.0.0.1", int(sys.argv[1]))) +def send(): + with open(sys.argv[2], "rb") as f: + while True: + b = f.read(65536) + if not b: + return + s.sendall(b) +threading.Thread(target=send, daemon=True).start() +got, paused = 0, False +with open(sys.argv[3], "wb") as out: + while True: + if not paused and got >= int(sys.argv[4]): + time.sleep(float(sys.argv[5])) + paused = True + d = s.recv(65536) + if not d: + break + out.write(d) + out.flush() + got += len(d) +' $echo_entry_port "$echo_payload" "$echo_received" $echo_pause_after \ + $echo_pause_secs 2>/dev/null & +echo_client_pid=$! + +counter=0 +got=0 +last=0 +stalled=0 +while [ "$got" -lt "$echo_payload_size" ] \ + && [ "$counter" -lt "$transfer_limit" ] +do + sleep 1 + counter=`expr $counter + 1` + got=`wc -c < "$echo_received" 2>/dev/null | tr -d ' '` + [ -z "$got" ] && got=0 + if [ "$got" -eq "$last" ] + then + stalled=`expr $stalled + 1` + [ "$stalled" -ge "$echo_stall_limit" ] && break + else + stalled=0 + last=$got + fi +done + +[ "$got" -eq "$echo_payload_size" ] \ + || do_fail "two-way forward stalled: sent $echo_payload_size, echoed $got" +cmp -s "$echo_payload" "$echo_received" \ + || do_fail "the echoed data does not match what was sent" + +echo "echoed $echo_payload_size bytes both ways while the reader stalled" + +# --- Phase 4: a new forward starts empty ------------------------------------ +# The echoserver serves one forward at a time, so wait for Phase 3's target +# to see its forward closed. +counter=0 +while kill -0 $echo_target_pid 2>/dev/null && [ "$counter" -lt 100 ]; do + sleep 0.1 + counter=`expr $counter + 1` +done +kill -0 $echo_target_pid 2>/dev/null \ + && do_fail "the echoserver didn't close the first forward" + +# Never reads the first connection, and closes its side once the client has +# stalled. Records what the second connection brings. +python3 -c ' +import os, socket, sys, time +s = socket.socket() +s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) +s.bind(("127.0.0.1", int(sys.argv[1]))) +s.listen(2) +open(sys.argv[2], "w").write("up") +a = s.accept()[0] +while not os.path.exists(sys.argv[3]): + time.sleep(0.1) +a.shutdown(socket.SHUT_WR) +time.sleep(0.5) +while a.recv(65536): + pass +open(sys.argv[4], "w").write("closed") +b = s.accept()[0] +b.settimeout(5) +got = b"" +try: + while len(got) < int(sys.argv[6]): + d = b.recv(65536) + if not d: + break + got += d +except socket.timeout: + pass +open(sys.argv[5], "wb").write(got) +' $reuse_target_port "$reuse_target_ready" "$reuse_stalled" "$reuse_closed" \ + "$reuse_received" `printf '%s' "$reuse_message" | wc -c` 2>/dev/null & +reuse_target_pid=$! + +counter=0 +while [ ! -s "$reuse_target_ready" ] && [ "$counter" -lt 30 ]; do + sleep 0.1 + counter=`expr $counter + 1` +done +[ -s "$reuse_target_ready" ] \ + || do_fail "couldn't listen on port $reuse_target_port, is it in use?" + +# Sends until a second passes with nothing taken, then opens the second +# connection once the first is closed. +python3 -c ' +import os, socket, sys, time +a = socket.create_connection(("127.0.0.1", int(sys.argv[1]))) +a.setblocking(False) +chunk = b"A" * 65536 +deadline = time.time() + 20 +idle = time.time() +while time.time() - idle < 1 and time.time() < deadline: + try: + if a.send(chunk) > 0: + idle = time.time() + except BlockingIOError: + time.sleep(0.05) +open(sys.argv[2], "w").write("stalled") +while not os.path.exists(sys.argv[3]): + time.sleep(0.1) +b = socket.create_connection(("127.0.0.1", int(sys.argv[1]))) +b.sendall(sys.argv[4].encode()) +time.sleep(60) +' $reuse_entry_port "$reuse_stalled" "$reuse_closed" "$reuse_message" \ + 2>/dev/null & +reuse_client_pid=$! + +counter=0 +while [ ! -f "$reuse_received" ] && [ "$counter" -lt "$reuse_limit" ]; do + sleep 1 + counter=`expr $counter + 1` +done +[ -f "$reuse_received" ] \ + || do_fail "the second forward never reached its target" +printf '%s' "$reuse_message" > "$reuse_expected" +cmp -s "$reuse_expected" "$reuse_received" \ + || do_fail "the second forward's target got `wc -c < "$reuse_received" \ + | tr -d ' '` bytes, expected only \"$reuse_message\"" + +echo "a new forward carried only its own bytes" do_cleanup exit 0