|
|
|
@ -205,6 +205,7 @@ static void flush_write_buf(struct tcp_conn* tc) { |
|
|
|
tc->write_buf = NULL; |
|
|
|
tc->write_buf = NULL; |
|
|
|
tc->write_len = 0; |
|
|
|
tc->write_len = 0; |
|
|
|
tc->write_offset = 0; |
|
|
|
tc->write_offset = 0; |
|
|
|
|
|
|
|
DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: flush_wbuf done fd=%d", (int)tc->sock); |
|
|
|
return; |
|
|
|
return; |
|
|
|
} |
|
|
|
} |
|
|
|
continue; |
|
|
|
continue; |
|
|
|
@ -220,19 +221,20 @@ static void flush_write_buf(struct tcp_conn* tc) { |
|
|
|
static void write_queue_fetch_cb(struct ll_queue* q, void* arg) { |
|
|
|
static void write_queue_fetch_cb(struct ll_queue* q, void* arg) { |
|
|
|
struct tcp_conn* tc = (struct tcp_conn*)arg; |
|
|
|
struct tcp_conn* tc = (struct tcp_conn*)arg; |
|
|
|
|
|
|
|
|
|
|
|
if (tc->error || tc->fin) return; |
|
|
|
if (tc->error || tc->fin) { DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: wq_fetch skip fd=%d error=%d fin=%d", (int)tc->sock, tc->error, tc->fin); return; } |
|
|
|
|
|
|
|
|
|
|
|
if (tc->write_buf) { |
|
|
|
if (tc->write_buf) { |
|
|
|
flush_write_buf(tc); |
|
|
|
flush_write_buf(tc); |
|
|
|
if (tc->write_buf) return; // остался остаток, ждём write_cb
|
|
|
|
if (tc->write_buf) return; // остался остаток, ждём write_cb
|
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
if (!tc->connected) return; // ждём connect
|
|
|
|
if (!tc->connected) { DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: wq_fetch skip fd=%d not connected", (int)tc->sock); return; } |
|
|
|
|
|
|
|
|
|
|
|
struct ll_entry* e = queue_data_get(q); |
|
|
|
struct ll_entry* e = queue_data_get(q); |
|
|
|
if (!e) { |
|
|
|
if (!e) { |
|
|
|
uasync_set_socket_write(tc->ua, tc->socket_id, 0); |
|
|
|
uasync_set_socket_write(tc->ua, tc->socket_id, 0); |
|
|
|
tc->write_monitor = 0; |
|
|
|
tc->write_monitor = 0; |
|
|
|
|
|
|
|
DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: wq_fetch POLLOUT OFF fd=%d (queue empty)", (int)tc->sock); |
|
|
|
queue_resume_callback(q); |
|
|
|
queue_resume_callback(q); |
|
|
|
if (tc->on_flushed) { |
|
|
|
if (tc->on_flushed) { |
|
|
|
void (*cb)(struct tcp_conn*, void*) = tc->on_flushed; |
|
|
|
void (*cb)(struct tcp_conn*, void*) = tc->on_flushed; |
|
|
|
@ -264,6 +266,7 @@ static void write_queue_fetch_cb(struct ll_queue* q, void* arg) { |
|
|
|
queue_entry_free(e); |
|
|
|
queue_entry_free(e); |
|
|
|
uasync_set_socket_write(tc->ua, tc->socket_id, 1); |
|
|
|
uasync_set_socket_write(tc->ua, tc->socket_id, 1); |
|
|
|
tc->write_monitor = 1; |
|
|
|
tc->write_monitor = 1; |
|
|
|
|
|
|
|
DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: wq_fetch POLLOUT ON fd=%d (partial send, wbuf=%zu bytes)", (int)tc->sock, tc->write_len); |
|
|
|
} else if (n < 0 && (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINTR)) { |
|
|
|
} else if (n < 0 && (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINTR)) { |
|
|
|
tc->write_buf = memory_pool_alloc(tc->data_pool); |
|
|
|
tc->write_buf = memory_pool_alloc(tc->data_pool); |
|
|
|
if (!tc->write_buf) { |
|
|
|
if (!tc->write_buf) { |
|
|
|
@ -280,6 +283,7 @@ static void write_queue_fetch_cb(struct ll_queue* q, void* arg) { |
|
|
|
queue_entry_free(e); |
|
|
|
queue_entry_free(e); |
|
|
|
uasync_set_socket_write(tc->ua, tc->socket_id, 1); |
|
|
|
uasync_set_socket_write(tc->ua, tc->socket_id, 1); |
|
|
|
tc->write_monitor = 1; |
|
|
|
tc->write_monitor = 1; |
|
|
|
|
|
|
|
DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: wq_fetch POLLOUT ON fd=%d (EAGAIN, wbuf=%zu bytes)", (int)tc->sock, tc->write_len); |
|
|
|
} else { |
|
|
|
} else { |
|
|
|
tc->error = 1; |
|
|
|
tc->error = 1; |
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: send error fd=%d errno=%d", (int)tc->sock, errno); |
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: send error fd=%d errno=%d", (int)tc->sock, errno); |
|
|
|
@ -302,6 +306,7 @@ static void write_cb(socket_t sock, void* arg) { |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_io: connect ok fd=%d", (int)tc->sock); |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_io: connect ok fd=%d", (int)tc->sock); |
|
|
|
} else { |
|
|
|
} else { |
|
|
|
tc->error = 1; |
|
|
|
tc->error = 1; |
|
|
|
|
|
|
|
DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: write_cb connect fail fd=%d err=%d wmon=%d (POLLOUT still active)", (int)tc->sock, err, tc->write_monitor); |
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: connect fail fd=%d err=%d", (int)tc->sock, err); |
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: connect fail fd=%d err=%d", (int)tc->sock, err); |
|
|
|
if (tc->on_error) tc->on_error(tc, err ? err : -1, tc->arg); |
|
|
|
if (tc->on_error) tc->on_error(tc, err ? err : -1, tc->arg); |
|
|
|
return; |
|
|
|
return; |
|
|
|
@ -309,7 +314,10 @@ static void write_cb(socket_t sock, void* arg) { |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
flush_write_buf(tc); |
|
|
|
flush_write_buf(tc); |
|
|
|
if (!tc->write_buf) queue_resume_callback(tc->write_queue); |
|
|
|
if (!tc->write_buf) { |
|
|
|
|
|
|
|
DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: write_cb flushed fd=%d wmon=%d (resume WQ cb)", (int)tc->sock, tc->write_monitor); |
|
|
|
|
|
|
|
queue_resume_callback(tc->write_queue); |
|
|
|
|
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
// ====================================================================
|
|
|
|
// ====================================================================
|
|
|
|
@ -321,6 +329,7 @@ static void error_cb(socket_t sock, void* arg) { |
|
|
|
struct tcp_conn* tc = (struct tcp_conn*)arg; |
|
|
|
struct tcp_conn* tc = (struct tcp_conn*)arg; |
|
|
|
if (!tc || tc->sock == SOCKET_INVALID) return; |
|
|
|
if (!tc || tc->sock == SOCKET_INVALID) return; |
|
|
|
tc->error = 1; |
|
|
|
tc->error = 1; |
|
|
|
|
|
|
|
DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: error_cb fd=%d connected=%d fin=%d wmon=%d", (int)tc->sock, tc->connected, tc->fin, tc->write_monitor); |
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: async error fd=%d", (int)tc->sock); |
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: async error fd=%d", (int)tc->sock); |
|
|
|
if (tc->on_error) tc->on_error(tc, -1, tc->arg); |
|
|
|
if (tc->on_error) tc->on_error(tc, -1, tc->arg); |
|
|
|
} |
|
|
|
} |
|
|
|
|