Proper flow-control handling

Patch by FRJ.
This commit is contained in:
Kim Woelders
2019-03-21 08:02:28 +01:00
parent 4d1d6978e4
commit 688fc53caf
+38 -17
View File
@@ -62,7 +62,9 @@ typedef struct {
struct fuse_pollhandle *ph; struct fuse_pollhandle *ph;
volatile char n_fds; volatile int rel_pending;
volatile int slave_suspended; /*waiting for application to read data */
volatile int master_suspended; /*server is not ready to receive data */
struct pollfd fds[3]; struct pollfd fds[3];
char net_buf[NET_BUF_SIZE]; char net_buf[NET_BUF_SIZE];
@@ -75,6 +77,7 @@ typedef struct {
int smcr_last; int smcr_last;
volatile int pollin; volatile int pollin;
volatile int pollout;
volatile int error; volatile int error;
int epipe; int epipe;
@@ -212,7 +215,7 @@ static void ttynvt_release(fuse_req_t req, struct fuse_file_info *info)
_log(LOG_INFO, "connection closed\n"); _log(LOG_INFO, "connection closed\n");
tty->n_fds = 0; tty->rel_pending = 1;
_fd_close(tty->fds[FD_NET].fd); _fd_close(tty->fds[FD_NET].fd);
_fd_close(tty->fds[FD_MASTER].fd); _fd_close(tty->fds[FD_MASTER].fd);
_fd_close(tty->fds[FD_SLAVE].fd); _fd_close(tty->fds[FD_SLAVE].fd);
@@ -278,13 +281,16 @@ static void *_read_net(void *arg)
pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL); pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL);
DBG2("%s: net_cnt=%d\n", __func__, tty->net_cnt); DBG2("%s: net_cnt=%d\n", __func__, tty->net_cnt);
while (tty->n_fds > 0) while (!tty->rel_pending)
{ {
tty->fds[FD_NET].revents = tty->fds[FD_MASTER].revents = tty->fds[FD_NET].revents = tty->fds[FD_MASTER].revents =
tty->fds[FD_SLAVE].revents = 0; tty->fds[FD_SLAVE].revents = 0;
res = poll(tty->fds, tty->n_fds, -1); tty->fds[FD_SLAVE].events = tty->slave_suspended ? POLLPRI : POLLIN;
DBG2("%s: res=%d N=%d events-N/M/S=%#x/%#x/%#x\n", __func__, tty->fds[FD_MASTER].events = tty->master_suspended ? POLLPRI : POLLIN;
res, tty->n_fds,
res = poll(tty->fds, 3, -1);
DBG2("%s: res=%d events-N/M/S=%#x/%#x/%#x\n", __func__,
res,
tty->fds[FD_NET].revents, tty->fds[FD_MASTER].revents, tty->fds[FD_NET].revents, tty->fds[FD_MASTER].revents,
tty->fds[FD_SLAVE].revents); tty->fds[FD_SLAVE].revents);
if (res < 0) if (res < 0)
@@ -293,7 +299,7 @@ static void *_read_net(void *arg)
continue; continue;
break; break;
} }
if (tty->n_fds == 0) if (tty->rel_pending)
break; break;
if (tty->fds[FD_NET].revents & POLLIN) if (tty->fds[FD_NET].revents & POLLIN)
@@ -330,7 +336,7 @@ static void *_read_net(void *arg)
break; break;
} }
if (tty->fds[FD_MASTER].revents & POLLIN) if ((tty->fds[FD_MASTER].revents & POLLIN) && !tty->master_suspended)
{ {
res = read(tty->fds[FD_MASTER].fd, buf, TMP_BUF_SIZE); res = read(tty->fds[FD_MASTER].fd, buf, TMP_BUF_SIZE);
if (res < 0) if (res < 0)
@@ -338,6 +344,11 @@ static void *_read_net(void *arg)
_log(LOG_ERR, "master read error: %m\n"); _log(LOG_ERR, "master read error: %m\n");
break; break;
} }
if (tty->pollout == 0)
{
tty->pollout = 1;
_notify(tty);
}
DBG2_BUF("TtyM in ", buf, res); DBG2_BUF("TtyM in ", buf, res);
telnet_tx(tty->tn, buf, res); telnet_tx(tty->tn, buf, res);
} }
@@ -345,7 +356,7 @@ static void *_read_net(void *arg)
if (tty->fds[FD_SLAVE].revents & POLLIN) if (tty->fds[FD_SLAVE].revents & POLLIN)
{ {
tty->pollin = 1; tty->pollin = 1;
tty->n_fds = 2; /* Disable polling of Slave */ tty->slave_suspended = 1; /* Disable polling of Slave */
_notify(tty); _notify(tty);
} }
} }
@@ -465,6 +476,8 @@ static void _modem_status_cb(void *cctx, int status)
if (status & TNS_STATE_CTS) if (status & TNS_STATE_CTS)
mcr |= TIOCM_CTS; mcr |= TIOCM_CTS;
tty->master_suspended = (mcr & TIOCM_CTS) ? 0 : 1;
pthread_mutex_lock(&tty->tty_lock); pthread_mutex_lock(&tty->tty_lock);
tty->smcr_last = tty->smcr; tty->smcr_last = tty->smcr;
tty->smcr = mcr; tty->smcr = mcr;
@@ -497,6 +510,8 @@ static void ttynvt_open(fuse_req_t req, struct fuse_file_info *info)
for (n = 0; n < 3; n++) for (n = 0; n < 3; n++)
tty->fds[n].fd = -1; tty->fds[n].fd = -1;
tty->pollout = 1;
tty->tn = telnet_ctx_init(tty, _srv_write, _srv_read, _modem_status_cb); tty->tn = telnet_ctx_init(tty, _srv_write, _srv_read, _modem_status_cb);
if (!tty->tn) if (!tty->tn)
{ {
@@ -565,8 +580,6 @@ static void ttynvt_open(fuse_req_t req, struct fuse_file_info *info)
goto open_err; goto open_err;
} }
tty->n_fds = 3; /* Initially poll Net, Master, and Slave */
info->fh = (uintptr_t) tty; info->fh = (uintptr_t) tty;
info->nonseekable = 1; info->nonseekable = 1;
info->direct_io = 1; info->direct_io = 1;
@@ -624,7 +637,7 @@ ttynvt_read(fuse_req_t req, size_t size, off_t off,
res = _is_interrupted ? -2 : read(tty->fds[FD_SLAVE].fd, buf, size); res = _is_interrupted ? -2 : read(tty->fds[FD_SLAVE].fd, buf, size);
_is_interrupted = 0; _is_interrupted = 0;
tty->ptid_read = 0; tty->ptid_read = 0;
tty->n_fds = 3; /* Enable polling of Slave */ tty->slave_suspended = 0; /* Enable polling of Slave */
fuse_req_interrupt_func(req, NULL, NULL); fuse_req_interrupt_func(req, NULL, NULL);
if (res < (int)size) if (res < (int)size)
{ {
@@ -680,6 +693,10 @@ ttynvt_write(fuse_req_t req, const char *data, size_t size, off_t off,
fuse_reply_err(req, errno); fuse_reply_err(req, errno);
return; return;
} }
else if (res != (int)size)
{
tty->pollout = 0;
}
fuse_reply_write(req, res); fuse_reply_write(req, res);
} }
@@ -870,7 +887,7 @@ ttynvt_ioctl(fuse_req_t req, int cmd, void *arg,
ioctl(tty->fds[FD_SLAVE].fd, cmd, tio); ioctl(tty->fds[FD_SLAVE].fd, cmd, tio);
if (cmd == TCSETSF) if (cmd == TCSETSF)
tty->n_fds = 3; /* Re-enable polling of Slave */ tty->slave_suspended = 0; /* Re-enable polling of Slave */
fuse_reply_ioctl(req, 0, 0, 0); fuse_reply_ioctl(req, 0, 0, 0);
} }
@@ -899,14 +916,16 @@ ttynvt_ioctl(fuse_req_t req, int cmd, void *arg,
default: default:
case TCIFLUSH: case TCIFLUSH:
byte = 1; byte = 1;
tty->n_fds = 3; /* Re-enable polling of Slave */ tty->slave_suspended = 0; /* Re-enable polling of Slave */
break; break;
case TCOFLUSH: case TCOFLUSH:
byte = 2; byte = 2;
tty->pollout = 1;
break; break;
case TCIOFLUSH: case TCIOFLUSH:
byte = 3; byte = 3;
tty->n_fds = 3; /* Re-enable polling of Slave */ tty->slave_suspended = 0; /* Re-enable polling of Slave */
tty->pollout = 1;
break; break;
} }
telnet_rfc2217_cfg(tty->tn, TNS_SET_PURGE, &byte, 1); telnet_rfc2217_cfg(tty->tn, TNS_SET_PURGE, &byte, 1);
@@ -1056,15 +1075,17 @@ ttynvt_poll(fuse_req_t req, struct fuse_file_info *info,
struct fuse_pollhandle *ph) struct fuse_pollhandle *ph)
{ {
ttynvt_t *tty = (ttynvt_t *) (uintptr_t) info->fh; ttynvt_t *tty = (ttynvt_t *) (uintptr_t) info->fh;
int revents; int revents = 0;
DBG2("%s: tty->pollin=%d tty->error/epipe=%d/%d\n", DBG2("%s: tty->pollin=%d tty->error/epipe=%d/%d\n",
__func__, tty->pollin, tty->error, tty->epipe); __func__, tty->pollin, tty->error, tty->epipe);
_update_notify(tty, ph); _update_notify(tty, ph);
if (tty->pollout)
{
revents = POLLOUT; revents = POLLOUT;
}
if (tty->pollin) if (tty->pollin)
{ {
revents |= POLLIN; revents |= POLLIN;