mirror of
https://github.com/cloudwu/skynet.git
synced 2026-07-25 12:43:09 +00:00
add log when socket send blocked by large data stream
This commit is contained in:
@@ -102,11 +102,17 @@ skynet_socket_poll() {
|
|||||||
|
|
||||||
int
|
int
|
||||||
skynet_socket_send(struct skynet_context *ctx, int id, void *buffer, int sz) {
|
skynet_socket_send(struct skynet_context *ctx, int id, void *buffer, int sz) {
|
||||||
int err = socket_server_send(SOCKET_SERVER, id, buffer, sz);
|
int64_t wsz = socket_server_send(SOCKET_SERVER, id, buffer, sz);
|
||||||
if (err < 0) {
|
if (wsz < 0) {
|
||||||
free(buffer);
|
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
|
int
|
||||||
|
|||||||
@@ -41,6 +41,7 @@ struct socket {
|
|||||||
int id;
|
int id;
|
||||||
int type;
|
int type;
|
||||||
int size;
|
int size;
|
||||||
|
int64_t wb_size;
|
||||||
uintptr_t opaque;
|
uintptr_t opaque;
|
||||||
struct write_buffer * head;
|
struct write_buffer * head;
|
||||||
struct write_buffer * tail;
|
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->fd = fd;
|
||||||
s->size = MIN_READ_BUFFER;
|
s->size = MIN_READ_BUFFER;
|
||||||
s->opaque = opaque;
|
s->opaque = opaque;
|
||||||
|
s->wb_size = 0;
|
||||||
assert(s->head == NULL);
|
assert(s->head == NULL);
|
||||||
assert(s->tail == NULL);
|
assert(s->tail == NULL);
|
||||||
return s;
|
return s;
|
||||||
@@ -345,6 +347,7 @@ send_buffer(struct socket_server *ss, struct socket *s, struct socket_message *r
|
|||||||
force_close(ss,s, result);
|
force_close(ss,s, result);
|
||||||
return SOCKET_CLOSE;
|
return SOCKET_CLOSE;
|
||||||
}
|
}
|
||||||
|
s->wb_size -= sz;
|
||||||
if (sz != tmp->sz) {
|
if (sz != tmp->sz) {
|
||||||
tmp->ptr += sz;
|
tmp->ptr += sz;
|
||||||
tmp->sz -= sz;
|
tmp->sz -= sz;
|
||||||
@@ -367,6 +370,24 @@ send_buffer(struct socket_server *ss, struct socket *s, struct socket_message *r
|
|||||||
return -1;
|
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
|
static int
|
||||||
send_socket(struct socket_server *ss, struct request_send * request, struct socket_message *result) {
|
send_socket(struct socket_server *ss, struct request_send * request, struct socket_message *result) {
|
||||||
int id = request->id;
|
int id = request->id;
|
||||||
@@ -396,25 +417,10 @@ send_socket(struct socket_server *ss, struct request_send * request, struct sock
|
|||||||
FREE(request->buffer);
|
FREE(request->buffer);
|
||||||
return -1;
|
return -1;
|
||||||
}
|
}
|
||||||
|
append_sendbuffer(s, request, n);
|
||||||
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;
|
|
||||||
|
|
||||||
sp_write(ss->event_fd, s->fd, s, true);
|
sp_write(ss->event_fd, s->fd, s, true);
|
||||||
} else {
|
} else {
|
||||||
struct write_buffer * buf = MALLOC(sizeof(*buf));
|
append_sendbuffer(s, request, 0);
|
||||||
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;
|
|
||||||
}
|
}
|
||||||
return -1;
|
return -1;
|
||||||
}
|
}
|
||||||
@@ -807,7 +813,7 @@ socket_server_block_connect(struct socket_server *ss, uintptr_t opaque, const ch
|
|||||||
}
|
}
|
||||||
|
|
||||||
// return -1 when error
|
// return -1 when error
|
||||||
int
|
int64_t
|
||||||
socket_server_send(struct socket_server *ss, int id, const void * buffer, int sz) {
|
socket_server_send(struct socket_server *ss, int id, const void * buffer, int sz) {
|
||||||
struct socket * s = &ss->slot[id % MAX_SOCKET];
|
struct socket * s = &ss->slot[id % MAX_SOCKET];
|
||||||
if (s->id != id || s->type == SOCKET_TYPE_INVALID) {
|
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;
|
request.u.send.buffer = (char *)buffer;
|
||||||
|
|
||||||
send_request(ss, &request, 'D', sizeof(request.u.send));
|
send_request(ss, &request, 'D', sizeof(request.u.send));
|
||||||
return 0;
|
return s->wb_size;
|
||||||
}
|
}
|
||||||
|
|
||||||
void
|
void
|
||||||
|
|||||||
@@ -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);
|
void socket_server_start(struct socket_server *, uintptr_t opaque, int id);
|
||||||
|
|
||||||
// return -1 when error
|
// 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
|
// ctrl command below returns id
|
||||||
int socket_server_listen(struct socket_server *, uintptr_t opaque, const char * addr, int port, int backlog);
|
int socket_server_listen(struct socket_server *, uintptr_t opaque, const char * addr, int port, int backlog);
|
||||||
|
|||||||
Reference in New Issue
Block a user