From 1230ea1a174586bff43209787b63b81bbc21c91d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=BA=91=E9=A3=8E?= Date: Mon, 13 Jan 2014 11:59:33 +0800 Subject: [PATCH] add log when socket send blocked by large data stream --- skynet-src/skynet_socket.c | 12 ++++++++--- skynet-src/socket_server.c | 44 ++++++++++++++++++++++---------------- skynet-src/socket_server.h | 2 +- 3 files changed, 35 insertions(+), 23 deletions(-) diff --git a/skynet-src/skynet_socket.c b/skynet-src/skynet_socket.c index c4c509d4..7ca48233 100644 --- a/skynet-src/skynet_socket.c +++ b/skynet-src/skynet_socket.c @@ -102,11 +102,17 @@ skynet_socket_poll() { int skynet_socket_send(struct skynet_context *ctx, int id, void *buffer, int sz) { - int err = socket_server_send(SOCKET_SERVER, id, buffer, sz); - if (err < 0) { + int64_t wsz = socket_server_send(SOCKET_SERVER, id, buffer, sz); + if (wsz < 0) { free(buffer); + return -1; + } else if (wsz > 1024 * 1024) { + int kb4 = wsz / 1024 / 4; + if (kb4 % 256 == 0) { + skynet_error(ctx, "%d Mb bytes on socket %d need to send out", (int)(wsz / (1024 * 1024)), id); + } } - return err; + return 0; } int diff --git a/skynet-src/socket_server.c b/skynet-src/socket_server.c index b83e05f8..8a76916b 100644 --- a/skynet-src/socket_server.c +++ b/skynet-src/socket_server.c @@ -41,6 +41,7 @@ struct socket { int id; int type; int size; + int64_t wb_size; uintptr_t opaque; struct write_buffer * head; struct write_buffer * tail; @@ -246,6 +247,7 @@ new_fd(struct socket_server *ss, int id, int fd, uintptr_t opaque, bool add) { s->fd = fd; s->size = MIN_READ_BUFFER; s->opaque = opaque; + s->wb_size = 0; assert(s->head == NULL); assert(s->tail == NULL); return s; @@ -345,6 +347,7 @@ send_buffer(struct socket_server *ss, struct socket *s, struct socket_message *r force_close(ss,s, result); return SOCKET_CLOSE; } + s->wb_size -= sz; if (sz != tmp->sz) { tmp->ptr += sz; tmp->sz -= sz; @@ -367,6 +370,24 @@ send_buffer(struct socket_server *ss, struct socket *s, struct socket_message *r return -1; } +static void +append_sendbuffer(struct socket *s, struct request_send * request, int n) { + struct write_buffer * buf = MALLOC(sizeof(*buf)); + buf->ptr = request->buffer+n; + buf->sz = request->sz - n; + buf->buffer = request->buffer; + buf->next = NULL; + s->wb_size += buf->sz; + if (s->head == NULL) { + s->head = s->tail = buf; + } else { + assert(s->tail != NULL); + assert(s->tail->next == NULL); + s->tail->next = buf; + s->tail = buf; + } +} + static int send_socket(struct socket_server *ss, struct request_send * request, struct socket_message *result) { int id = request->id; @@ -396,25 +417,10 @@ send_socket(struct socket_server *ss, struct request_send * request, struct sock FREE(request->buffer); return -1; } - - struct write_buffer * buf = MALLOC(sizeof(*buf)); - buf->next = NULL; - buf->ptr = request->buffer+n; - buf->sz = request->sz - n; - buf->buffer = request->buffer; - s->head = s->tail = buf; - + append_sendbuffer(s, request, n); sp_write(ss->event_fd, s->fd, s, true); } else { - struct write_buffer * buf = MALLOC(sizeof(*buf)); - buf->ptr = request->buffer; - buf->buffer = request->buffer; - buf->sz = request->sz; - assert(s->tail != NULL); - assert(s->tail->next == NULL); - buf->next = s->tail->next; - s->tail->next = buf; - s->tail = buf; + append_sendbuffer(s, request, 0); } return -1; } @@ -807,7 +813,7 @@ socket_server_block_connect(struct socket_server *ss, uintptr_t opaque, const ch } // return -1 when error -int +int64_t socket_server_send(struct socket_server *ss, int id, const void * buffer, int sz) { struct socket * s = &ss->slot[id % MAX_SOCKET]; if (s->id != id || s->type == SOCKET_TYPE_INVALID) { @@ -821,7 +827,7 @@ socket_server_send(struct socket_server *ss, int id, const void * buffer, int sz request.u.send.buffer = (char *)buffer; send_request(ss, &request, 'D', sizeof(request.u.send)); - return 0; + return s->wb_size; } void diff --git a/skynet-src/socket_server.h b/skynet-src/socket_server.h index f77dc53c..41351331 100644 --- a/skynet-src/socket_server.h +++ b/skynet-src/socket_server.h @@ -28,7 +28,7 @@ void socket_server_close(struct socket_server *, uintptr_t opaque, int id); void socket_server_start(struct socket_server *, uintptr_t opaque, int id); // return -1 when error -int socket_server_send(struct socket_server *, int id, const void * buffer, int sz); +int64_t socket_server_send(struct socket_server *, int id, const void * buffer, int sz); // ctrl command below returns id int socket_server_listen(struct socket_server *, uintptr_t opaque, const char * addr, int port, int backlog);