recv/send directly instead of switching over to a different task

This commit is contained in:
Witold Kręcicki
2018-08-20 16:02:49 +02:00
parent 3ab9a4f6bd
commit fa2c7fe216
3 changed files with 36 additions and 273 deletions
+5 -5
View File
@@ -213,7 +213,7 @@ task_finished(isc__task_t *task) {
* any idle worker threads so they
* can exit.
*/
for (int i=0; i<manager->workers; i++) {
for (unsigned i=0; i<manager->workers; i++) {
LOCK(&manager->locks[i]);
BROADCAST(&manager->work_available[i]);
UNLOCK(&manager->locks[i]);
@@ -1159,11 +1159,11 @@ dispatch(isc__taskmgr_t *manager, int c) {
if (manager->tasks_running == 0 && empty_readyq(manager, c)) {
if (manager->mode != isc_taskmgrmode_normal) {
manager->mode = isc_taskmgrmode_normal;
for (int i=0; i<manager->workers; i++) {
if (i != c)
for (unsigned i=0; i<manager->workers; i++) {
if (i != (unsigned)c)
LOCK(&manager->locks[i]);
BROADCAST(&manager->work_available[i]);
if (i != c)
if (i != (unsigned)c)
UNLOCK(&manager->locks[i]);
}
}
@@ -1542,7 +1542,7 @@ isc_task_endexclusive(isc_task_t *task0) {
REQUIRE(manager->exclusive_requested);
manager->exclusive_requested = false;
UNLOCK(&manager->lock);
for (int i=0; i < manager->workers; i++) {
for (unsigned i=0; i < manager->workers; i++) {
LOCK(&manager->locks[i]);
BROADCAST(&manager->work_available[i]);
UNLOCK(&manager->locks[i]);
+30 -267
View File
@@ -361,9 +361,7 @@ struct isc__socket {
isc_sockaddr_t peer_address; /* remote address */
unsigned int pending_recv : 1,
pending_send : 1,
pending_accept : 1,
unsigned int pending_accept : 1,
listener : 1, /* listener socket */
connected : 1,
connecting : 1, /* connect pending */
@@ -468,10 +466,6 @@ static isc_result_t allocate_socket(isc__socketmgr_t *, isc_sockettype_t,
static void destroy(isc__socket_t **);
static void internal_accept(isc_task_t *, isc_event_t *);
static void internal_connect(isc_task_t *, isc_event_t *);
static void internal_recv(isc_task_t *, isc_event_t *);
static void internal_send(isc_task_t *, isc_event_t *);
static void internal_fdwatch_write(isc_task_t *, isc_event_t *);
static void internal_fdwatch_read(isc_task_t *, isc_event_t *);
static void process_cmsg(isc__socket_t *, struct msghdr *, isc_socketevent_t *);
static void build_msghdr_send(isc__socket_t *, char *, isc_socketevent_t *,
struct msghdr *, struct iovec *, size_t *);
@@ -756,7 +750,6 @@ watch_fd(isc__socketmgr_t *manager, int fd, int msg) {
UNEXPECTED_ERROR(__FILE__, __LINE__,
"epoll_ctl(ADD/MOD) returned "
"EEXIST for fd %d", fd);
printf("CTL D failed %d %d %d %d %d %u %u %d\n", ret, errno, fd, manager->epoll_fds[fd % NUM_EPOLLS], op, oldevents, event.events, fd == manager->pipe_fds[0]);
result = isc__errno2result(errno);
}
return (result);
@@ -2148,8 +2141,6 @@ allocate_socket(isc__socketmgr_t *manager, isc_sockettype_t type,
ISC_LIST_INIT(sock->send_list);
ISC_LIST_INIT(sock->accept_list);
ISC_LIST_INIT(sock->connect_list);
sock->pending_recv = 0;
sock->pending_send = 0;
sock->pending_accept = 0;
sock->listener = 0;
sock->connected = 0;
@@ -2209,8 +2200,6 @@ free_socket(isc__socket_t **socketp) {
INSIST(VALID_SOCKET(sock));
INSIST(sock->references == 0);
INSIST(!sock->connecting);
INSIST(!sock->pending_recv);
INSIST(!sock->pending_send);
INSIST(!sock->pending_accept);
INSIST(ISC_LIST_EMPTY(sock->recv_list));
INSIST(ISC_LIST_EMPTY(sock->send_list));
@@ -2984,16 +2973,14 @@ isc_socket_fdwatchpoke(isc_socket_t *sock0, int flags)
if ((flags & (ISC_SOCKFDWATCH_READ | ISC_SOCKFDWATCH_WRITE)) != 0) {
LOCK(&sock->recvlock);
bool doit = (((flags & ISC_SOCKFDWATCH_READ) != 0) &&
!sock->pending_recv);
bool doit = (((flags & ISC_SOCKFDWATCH_READ) != 0));
UNLOCK(&sock->recvlock);
if (doit)
select_poke(sock->manager, sock->fd,
SELECT_POKE_READ);
LOCK(&sock->sendlock);
doit = (((flags & ISC_SOCKFDWATCH_WRITE) != 0) &&
!sock->pending_send);
doit = (((flags & ISC_SOCKFDWATCH_WRITE) != 0));
UNLOCK(&sock->sendlock);
if (doit)
select_poke(sock->manager, sock->fd,
@@ -3070,8 +3057,6 @@ isc_socket_close(isc_socket_t *sock0) {
REQUIRE(sock->fd >= 0 && sock->fd < (int)sock->manager->maxsocks);
INSIST(!sock->connecting);
INSIST(!sock->pending_recv);
INSIST(!sock->pending_send);
INSIST(!sock->pending_accept);
INSIST(ISC_LIST_EMPTY(sock->recv_list));
INSIST(ISC_LIST_EMPTY(sock->send_list));
@@ -3098,83 +3083,6 @@ isc_socket_close(isc_socket_t *sock0) {
return (ISC_R_SUCCESS);
}
/*
* I/O is possible on a given socket. Schedule an event to this task that
* will call an internal function to do the I/O. This will charge the
* task with the I/O operation and let our select loop handler get back
* to doing something real as fast as possible.
*
* The socket and manager must be locked before calling this function.
*/
static void
dispatch_recv(isc__socket_t *sock) {
intev_t *iev;
isc_socketevent_t *ev;
isc_task_t *sender;
INSIST(!sock->pending_recv);
if (sock->type != isc_sockettype_fdwatch) {
ev = ISC_LIST_HEAD(sock->recv_list);
if (ev == NULL)
return;
socket_log(sock, NULL, EVENT, NULL, 0, 0,
"dispatch_recv: event %p -> task %p",
ev, ev->ev_sender);
sender = ev->ev_sender;
} else {
sender = sock->fdwatchtask;
}
sock->pending_recv = 1;
iev = &sock->readable_ev;
sock->references++;
iev->ev_sender = sock;
if (sock->type == isc_sockettype_fdwatch)
iev->ev_action = internal_fdwatch_read;
else
iev->ev_action = internal_recv;
iev->ev_arg = sock;
UNLOCK(&sock->recvlock);
isc_task_send(sender, (isc_event_t **)&iev);
LOCK(&sock->recvlock);
}
static void
dispatch_send(isc__socket_t *sock) {
intev_t *iev;
isc_socketevent_t *ev;
isc_task_t *sender;
INSIST(!sock->pending_send);
if (sock->type != isc_sockettype_fdwatch) {
ev = ISC_LIST_HEAD(sock->send_list);
if (ev == NULL)
return;
socket_log(sock, NULL, EVENT, NULL, 0, 0,
"dispatch_send: event %p -> task %p",
ev, ev->ev_sender);
sender = ev->ev_sender;
} else {
sender = sock->fdwatchtask;
}
sock->pending_send = 1;
iev = &sock->writable_ev;
sock->references++;
iev->ev_sender = sock;
if (sock->type == isc_sockettype_fdwatch)
iev->ev_action = internal_fdwatch_write;
else
iev->ev_action = internal_send;
iev->ev_arg = sock;
isc_task_send(sender, (isc_event_t **)&iev);
}
/*
* Dispatch an internal accept event.
*/
@@ -3570,34 +3478,14 @@ internal_accept(isc_task_t *me, isc_event_t *ev) {
}
static void
internal_recv(isc_task_t *me, isc_event_t *ev) {
internal_recv(isc__socket_t *sock) {
isc_socketevent_t *dev;
isc__socket_t *sock;
bool locked = true;
INSIST(ev->ev_type == ISC_SOCKEVENT_INTR);
sock = ev->ev_sender;
INSIST(VALID_SOCKET(sock));
socket_log(sock, NULL, IOEVENT,
isc_msgcat, ISC_MSGSET_SOCKET, ISC_MSG_INTERNALRECV,
"internal_recv: task %p got event %p", me, ev);
LOCK(&sock->sendlock);
LOCK(&sock->recvlock);
INSIST(sock->pending_recv == 1);
sock->pending_recv = 0;
INSIST(sock->references > 0);
sock->references--; /* the internal event is done with this socket */
if (sock->references == 0) {
UNLOCK(&sock->recvlock);
destroy(&sock);
return;
}
UNLOCK(&sock->recvlock);
UNLOCK(&sock->sendlock);
LOCK(&sock->recvlock);
/*
* Try to do as much I/O as possible on this socket. There are no
@@ -3635,39 +3523,20 @@ internal_recv(isc_task_t *me, isc_event_t *ev) {
}
poke:
if (!ISC_LIST_EMPTY(sock->recv_list)) {
if (locked)
UNLOCK(&sock->recvlock);
select_poke(sock->manager, sock->fd, SELECT_POKE_READ);
} else {
if (locked)
UNLOCK(&sock->recvlock);
}
if (!ISC_LIST_EMPTY(sock->recv_list))
watch_fd(sock->manager, sock->fd, SELECT_POKE_READ);
UNLOCK(&sock->recvlock);
}
static void
internal_send(isc_task_t *me, isc_event_t *ev) {
internal_send(isc__socket_t *sock) {
isc_socketevent_t *dev;
isc__socket_t *sock;
INSIST(ev->ev_type == ISC_SOCKEVENT_INTW);
/*
* Find out what socket this is and lock it.
*/
sock = (isc__socket_t *)ev->ev_sender;
INSIST(VALID_SOCKET(sock));
LOCK(&sock->sendlock);
socket_log(sock, NULL, IOEVENT,
isc_msgcat, ISC_MSGSET_SOCKET, ISC_MSG_INTERNALSEND,
"internal_send: task %p got event %p", me, ev);
INSIST(sock->pending_send == 1);
sock->pending_send = 0;
INSIST(sock->references > 0);
sock->references--; /* the internal event is done with this socket */
if (sock->references == 0) {
UNLOCK(&sock->sendlock);
destroy(&sock);
@@ -3695,94 +3564,11 @@ internal_send(isc_task_t *me, isc_event_t *ev) {
poke:
if (!ISC_LIST_EMPTY(sock->send_list))
select_poke(sock->manager, sock->fd, SELECT_POKE_WRITE);
watch_fd(sock->manager, sock->fd, SELECT_POKE_WRITE);
UNLOCK(&sock->sendlock);
}
static void
internal_fdwatch_write(isc_task_t *me, isc_event_t *ev) {
isc__socket_t *sock;
int more_data;
INSIST(ev->ev_type == ISC_SOCKEVENT_INTW);
/*
* Find out what socket this is and lock it.
*/
sock = (isc__socket_t *)ev->ev_sender;
INSIST(VALID_SOCKET(sock));
LOCK(&sock->sendlock);
socket_log(sock, NULL, IOEVENT,
isc_msgcat, ISC_MSGSET_SOCKET, ISC_MSG_INTERNALSEND,
"internal_fdwatch_write: task %p got event %p", me, ev);
INSIST(sock->pending_send == 1);
UNLOCK(&sock->sendlock);
more_data = (sock->fdwatchcb)(me, (isc_socket_t *)sock,
sock->fdwatcharg, ISC_SOCKFDWATCH_WRITE);
LOCK(&sock->sendlock);
sock->pending_send = 0;
INSIST(sock->references > 0);
sock->references--; /* the internal event is done with this socket */
if (sock->references == 0) {
UNLOCK(&sock->sendlock);
destroy(&sock);
return;
}
if (more_data)
select_poke(sock->manager, sock->fd, SELECT_POKE_WRITE);
UNLOCK(&sock->sendlock);
}
static void
internal_fdwatch_read(isc_task_t *me, isc_event_t *ev) {
isc__socket_t *sock;
int more_data;
INSIST(ev->ev_type == ISC_SOCKEVENT_INTR);
/*
* Find out what socket this is and lock it.
*/
sock = (isc__socket_t *)ev->ev_sender;
INSIST(VALID_SOCKET(sock));
LOCK(&sock->recvlock);
socket_log(sock, NULL, IOEVENT,
isc_msgcat, ISC_MSGSET_SOCKET, ISC_MSG_INTERNALRECV,
"internal_fdwatch_read: task %p got event %p", me, ev);
INSIST(sock->pending_recv == 1);
UNLOCK(&sock->recvlock);
more_data = (sock->fdwatchcb)(me, (isc_socket_t *)sock,
sock->fdwatcharg, ISC_SOCKFDWATCH_READ);
LOCK(&sock->recvlock);
sock->pending_recv = 0;
INSIST(sock->references > 0);
sock->references--; /* the internal event is done with this socket */
if (sock->references == 0) {
UNLOCK(&sock->recvlock);
destroy(&sock);
return;
}
UNLOCK(&sock->recvlock);
if (more_data)
select_poke(sock->manager, sock->fd, SELECT_POKE_READ);
}
/*
* Process read/writes on each fd here. Avoid locking
* and unlocking twice if both reads and writes are possible.
@@ -3815,13 +3601,13 @@ process_fd(isc__socketmgr_t *manager, int fd, bool readable,
unwatch_read = true;
goto check_write;
}
unlock_sock_recv = true;
LOCK(&sock->recvlock);
if (!SOCK_DEAD(sock)) {
if (sock->listener)
if (sock->listener) {
unlock_sock_recv = true;
LOCK(&sock->recvlock);
dispatch_accept(sock);
else
dispatch_recv(sock);
} else
internal_recv(sock);
}
unwatch_read = true;
}
@@ -3831,14 +3617,13 @@ check_write:
unwatch_write = true;
goto unlock_fd;
}
unlock_sock_write = true;
LOCK(&sock->sendlock);
if (!SOCK_DEAD(sock)) {
if (sock->connecting)
if (sock->connecting) {
unlock_sock_write = true;
LOCK(&sock->sendlock);
dispatch_connect(sock);
else
dispatch_send(sock);
} else
internal_send(sock);
}
unwatch_write = true;
}
@@ -4287,7 +4072,6 @@ setup_watcher(isc_mem_t *mctx, isc__socketmgr_t *manager) {
UNEXPECTED_ERROR(__FILE__, __LINE__,
"epoll_ctl(ADD/MOD) returned "
"EEXIST for fd %d", fd);
printf("CTL failed %d %d %d %d %d %u %u %d\n", ret, errno, fd, manager->epoll_fds[fd % NUM_EPOLLS], op, oldevents, event.events, fd == manager->pipe_fds[0]);
result = isc__errno2result(errno);
return (result);
}
@@ -4403,7 +4187,7 @@ setup_watcher(isc_mem_t *mctx, isc__socketmgr_t *manager) {
static void
cleanup_watcher(isc_mem_t *mctx, isc__socketmgr_t *manager) {
isc_result_t result;
isc_result_t result = ISC_R_SUCCESS;
// FOO result = unwatch_fd(manager, manager->pipe_fds[0], SELECT_POKE_READ);
if (result != ISC_R_SUCCESS) {
@@ -4529,7 +4313,7 @@ isc_socketmgr_create2(isc_mem_t *mctx, isc_socketmgr_t **managerp,
* Create the special fds that will be used to wake up the
* select/poll loop when something internal needs to be done.
*/
for (int i=0; i < NUM_EPOLLS; i++) {
for (i=0; i < NUM_EPOLLS; i++) {
if (pipe(manager->pipe_fds[i]) != 0) {
isc__strerror(errno, strbuf, sizeof(strbuf));
UNEXPECTED_ERROR(__FILE__, __LINE__,
@@ -4556,7 +4340,7 @@ isc_socketmgr_create2(isc_mem_t *mctx, isc_socketmgr_t **managerp,
* Start up the select/poll thread.
*/
isc_mem_attach(mctx, &manager->mctx);
for (int i=0; i<NUM_EPOLLS; i++) {
for (i=0; i<NUM_EPOLLS; i++) {
isc__mgrplus_t *m = isc_mem_get(mctx, sizeof(isc__mgrplus_t));
m->manager = manager;
m->id = i;
@@ -4586,7 +4370,7 @@ isc_socketmgr_create2(isc_mem_t *mctx, isc_socketmgr_t **managerp,
return (ISC_R_SUCCESS);
cleanup:
for (int i=0; i< NUM_EPOLLS; i++) {
for (i=0; i< NUM_EPOLLS; i++) {
(void)close(manager->pipe_fds[i][0]);
(void)close(manager->pipe_fds[i][1]);
}
@@ -4683,13 +4467,13 @@ isc_socketmgr_destroy(isc_socketmgr_t **managerp) {
* half of the pipe, which will send EOF to the read half.
* This is currently a no-op in the non-threaded case.
*/
for (int i=0; i<NUM_EPOLLS; i++)
for (i=0; i<NUM_EPOLLS; i++)
select_poke(manager, i, SELECT_POKE_SHUTDOWN);
/*
* Wait for thread to exit.
*/
for (int i=0; i<NUM_EPOLLS; i++) {
for (i=0; i<NUM_EPOLLS; i++) {
if (isc_thread_join(manager->watcher[i], NULL) != ISC_R_SUCCESS)
UNEXPECTED_ERROR(__FILE__, __LINE__,
"isc_thread_join() %s",
@@ -4702,7 +4486,7 @@ isc_socketmgr_destroy(isc_socketmgr_t **managerp) {
*/
cleanup_watcher(manager->mctx, manager);
for (int i=0; i<NUM_EPOLLS; i++) {
for (i=0; i<NUM_EPOLLS; i++) {
(void)close(manager->pipe_fds[i][0]);
(void)close(manager->pipe_fds[i][1]);
}
@@ -4786,7 +4570,7 @@ socket_recv(isc__socket_t *sock, isc_socketevent_t *dev, isc_task_t *task,
* Enqueue the request. If the socket was previously not being
* watched, poke the watcher to start paying attention to it.
*/
poke = (ISC_LIST_EMPTY(sock->recv_list) && !sock->pending_recv);
poke = (ISC_LIST_EMPTY(sock->recv_list));
ISC_LIST_ENQUEUE(sock->recv_list, dev, ev_link);
@@ -4995,8 +4779,7 @@ socket_send(isc__socket_t *sock, isc_socketevent_t *dev, isc_task_t *task,
* not being watched, poke the watcher to start
* paying attention to it.
*/
if (ISC_LIST_EMPTY(sock->send_list) &&
!sock->pending_send)
if (ISC_LIST_EMPTY(sock->send_list))
select_poke(sock->manager, sock->fd,
SELECT_POKE_WRITE);
ISC_LIST_ENQUEUE(sock->send_list, dev, ev_link);
@@ -6312,14 +6095,6 @@ isc_socketmgr_renderxml(isc_socketmgr_t *mgr0, xmlTextWriterPtr writer) {
}
TRY0(xmlTextWriterStartElement(writer, ISC_XMLCHAR "states"));
if (sock->pending_recv)
TRY0(xmlTextWriterWriteElement(writer,
ISC_XMLCHAR "state",
ISC_XMLCHAR "pending-receive"));
if (sock->pending_send)
TRY0(xmlTextWriterWriteElement(writer,
ISC_XMLCHAR "state",
ISC_XMLCHAR "pending-send"));
if (sock->pending_accept)
TRY0(xmlTextWriterWriteElement(writer,
ISC_XMLCHAR "state",
@@ -6431,18 +6206,6 @@ isc_socketmgr_renderjson(isc_socketmgr_t *mgr0, json_object *stats) {
CHECKMEM(states);
json_object_object_add(entry, "states", states);
if (sock->pending_recv) {
obj = json_object_new_string("pending-receive");
CHECKMEM(obj);
json_object_array_add(states, obj);
}
if (sock->pending_send) {
obj = json_object_new_string("pending-send");
CHECKMEM(obj);
json_object_array_add(states, obj);
}
if (sock->pending_accept) {
obj = json_object_new_string("pending-accept");
CHECKMEM(obj);
+1 -1
View File
@@ -60,7 +60,7 @@ static volatile HANDLE _mutex = NULL;
#else /* defined(_WIN32) || defined(_WIN64) */
#include <pthread.h>
static pthread_mutex_t _mutex = PTHREAD_MUTEX_INITIALIZER;
//static pthread_mutex_t _mutex = PTHREAD_MUTEX_INITIALIZER;
#define _LOCK() pthread_mutex_lock(&_mutex)
#define _UNLOCK() pthread_mutex_unlock(&_mutex)
#endif /* defined(_WIN32) || defined(_WIN64) */