|
|
|
@ -31,8 +31,13 @@ static struct remote_proxy_ctx* g_rp_ctx = NULL; |
|
|
|
static void rp_sock_read_cb(socket_t sock, void* arg); |
|
|
|
static void rp_sock_read_cb(socket_t sock, void* arg); |
|
|
|
static void rp_sock_write_cb(socket_t sock, void* arg); |
|
|
|
static void rp_sock_write_cb(socket_t sock, void* arg); |
|
|
|
static void rp_sock_error_cb(socket_t sock, void* arg); |
|
|
|
static void rp_sock_error_cb(socket_t sock, void* arg); |
|
|
|
|
|
|
|
static void rp_sock_retry_cb(void* arg); |
|
|
|
|
|
|
|
static void rp_sock_pause_cb(void* arg); |
|
|
|
|
|
|
|
static void rp_sock_send(struct remote_proxy_conn* rc, const uint8_t* data, size_t len); |
|
|
|
void rp_conn_free(struct remote_proxy_conn* rc); |
|
|
|
void rp_conn_free(struct remote_proxy_conn* rc); |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
#define RP_NORMALIZER_Q_THRESHOLD 64 |
|
|
|
|
|
|
|
|
|
|
|
static int rp_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, |
|
|
|
static int rp_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, |
|
|
|
uint32_t sid, uint16_t seq, const uint8_t* data, size_t len) { |
|
|
|
uint32_t sid, uint16_t seq, const uint8_t* data, size_t len) { |
|
|
|
struct ll_entry* e = queue_entry_new(0); |
|
|
|
struct ll_entry* e = queue_entry_new(0); |
|
|
|
@ -62,6 +67,89 @@ static int rp_conn_total(struct remote_proxy_conn* rc) { |
|
|
|
return n; |
|
|
|
return n; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// ====================================================================
|
|
|
|
|
|
|
|
// Неблокирующий send с EAGAIN и таймером повтора
|
|
|
|
|
|
|
|
// ====================================================================
|
|
|
|
|
|
|
|
static int rp_sock_try_send(struct remote_proxy_conn* rc) { |
|
|
|
|
|
|
|
if (!rc || rc->sock == SOCKET_INVALID || !rc->out_buf) return 0; |
|
|
|
|
|
|
|
while (rc->out_off < rc->out_len) { |
|
|
|
|
|
|
|
ssize_t n = send(rc->sock, rc->out_buf + rc->out_off, rc->out_len - rc->out_off, MSG_NOSIGNAL); |
|
|
|
|
|
|
|
if (n > 0) { |
|
|
|
|
|
|
|
rc->out_off += n; |
|
|
|
|
|
|
|
if (rc->out_off >= rc->out_len) { |
|
|
|
|
|
|
|
u_free(rc->out_buf); rc->out_buf = NULL; rc->out_len = rc->out_off = 0; |
|
|
|
|
|
|
|
return 0; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
continue; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
if (n < 0 && (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINTR)) return 1; |
|
|
|
|
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:SEND_ERR fd=%d sid=%08x err=%s", (int)rc->sock, rc->stream_id, strerror(errno)); |
|
|
|
|
|
|
|
return -1; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
return 0; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
static void rp_sock_retry_cb(void* arg) { |
|
|
|
|
|
|
|
struct remote_proxy_conn* rc = (struct remote_proxy_conn*)arg; |
|
|
|
|
|
|
|
if (!rc || !rc->out_buf) { rc->out_timer = NULL; return; } |
|
|
|
|
|
|
|
int ret = rp_sock_try_send(rc); |
|
|
|
|
|
|
|
if (ret == 0) { |
|
|
|
|
|
|
|
rc->out_timer = NULL; rc->out_backoff = 0; |
|
|
|
|
|
|
|
} else if (ret == 1) { |
|
|
|
|
|
|
|
rc->out_backoff = rc->out_backoff < 20 ? 20 : rc->out_backoff * 2; |
|
|
|
|
|
|
|
if (rc->out_backoff > 5000) rc->out_backoff = 5000; |
|
|
|
|
|
|
|
rc->out_timer = uasync_set_timeout(rc->ua, rc->out_backoff, rc, rp_sock_retry_cb, "rp_send"); |
|
|
|
|
|
|
|
} else { |
|
|
|
|
|
|
|
u_free(rc->out_buf); rc->out_buf = NULL; rc->out_len = rc->out_off = 0; |
|
|
|
|
|
|
|
rc->out_timer = NULL; rc->out_backoff = 0; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
static void rp_sock_send(struct remote_proxy_conn* rc, const uint8_t* data, size_t len) { |
|
|
|
|
|
|
|
if (!rc || !data || len == 0) return; |
|
|
|
|
|
|
|
if (rc->out_buf) { |
|
|
|
|
|
|
|
size_t new_len = rc->out_len + len; |
|
|
|
|
|
|
|
uint8_t* new_buf = u_realloc(rc->out_buf, new_len); |
|
|
|
|
|
|
|
if (!new_buf) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:SEND realloc(%zu) failed sid=%08x", new_len, rc->stream_id); return; } |
|
|
|
|
|
|
|
memcpy(new_buf + rc->out_len, data, len); |
|
|
|
|
|
|
|
rc->out_buf = new_buf; rc->out_len = new_len; |
|
|
|
|
|
|
|
return; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
rc->out_buf = u_malloc(len); |
|
|
|
|
|
|
|
if (!rc->out_buf) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:SEND malloc(%zu) failed sid=%08x", len, rc->stream_id); return; } |
|
|
|
|
|
|
|
memcpy(rc->out_buf, data, len); rc->out_len = len; rc->out_off = 0; |
|
|
|
|
|
|
|
int ret = rp_sock_try_send(rc); |
|
|
|
|
|
|
|
if (ret == 0) return; |
|
|
|
|
|
|
|
if (ret == 1) { |
|
|
|
|
|
|
|
rc->out_backoff = 10; |
|
|
|
|
|
|
|
rc->out_timer = uasync_set_timeout(rc->ua, rc->out_backoff, rc, rp_sock_retry_cb, "rp_send"); |
|
|
|
|
|
|
|
} else { |
|
|
|
|
|
|
|
u_free(rc->out_buf); rc->out_buf = NULL; rc->out_len = rc->out_off = 0; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
static void rp_sock_pause_cb(void* arg) { |
|
|
|
|
|
|
|
struct remote_proxy_conn* rc = (struct remote_proxy_conn*)arg; |
|
|
|
|
|
|
|
if (!rc || rc->sock == SOCKET_INVALID) { rc->pause_timer = NULL; return; } |
|
|
|
|
|
|
|
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; |
|
|
|
|
|
|
|
int qcnt = -1; |
|
|
|
|
|
|
|
if (inst) qcnt = etcp_router_input_q_count(inst, rc->peer_node_id); |
|
|
|
|
|
|
|
if (qcnt > RP_NORMALIZER_Q_THRESHOLD) { |
|
|
|
|
|
|
|
rc->pause_backoff = rc->pause_backoff * 2; |
|
|
|
|
|
|
|
if (rc->pause_backoff > 5000) rc->pause_backoff = 5000; |
|
|
|
|
|
|
|
rc->pause_timer = uasync_set_timeout(rc->ua, rc->pause_backoff, rc, rp_sock_pause_cb, "rp_pause"); |
|
|
|
|
|
|
|
return; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
if (rc->pause_buf && rc->pause_len > 0) { |
|
|
|
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:PAUSE_FLUSH fd=%d sid=%08x len=%zu qcnt=%d", (int)rc->sock, rc->stream_id, rc->pause_len, qcnt); |
|
|
|
|
|
|
|
rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, rc->send_seq, rc->pause_buf, rc->pause_len); |
|
|
|
|
|
|
|
rc->send_seq++; |
|
|
|
|
|
|
|
u_free(rc->pause_buf); rc->pause_buf = NULL; rc->pause_len = 0; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
rc->read_id = uasync_add_socket_t(rc->ua, rc->sock, rp_sock_read_cb, rp_sock_write_cb, rp_sock_error_cb, rc); |
|
|
|
|
|
|
|
rc->pause_timer = NULL; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
// ====================================================================
|
|
|
|
// ====================================================================
|
|
|
|
// Socket callbacks
|
|
|
|
// Socket callbacks
|
|
|
|
// ====================================================================
|
|
|
|
// ====================================================================
|
|
|
|
@ -73,11 +161,20 @@ static void rp_sock_read_cb(socket_t sock, void* arg) { |
|
|
|
if (rc->cli_closed || rc->error) return; |
|
|
|
if (rc->cli_closed || rc->error) return; |
|
|
|
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; |
|
|
|
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; |
|
|
|
if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:RECV fd=%d sid=%08x — inst is NULL, drop", (int)rc->sock, rc->stream_id); return; } |
|
|
|
if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:RECV fd=%d sid=%08x — inst is NULL, drop", (int)rc->sock, rc->stream_id); return; } |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:RECV fd=%d sid=%08x len=%zd total=%d", (int)rc->sock, rc->stream_id, n, rp_conn_total(rc)); |
|
|
|
int qcnt = etcp_router_input_q_count(inst, rc->peer_node_id); |
|
|
|
if (!rc->error) { |
|
|
|
if (qcnt > RP_NORMALIZER_Q_THRESHOLD) { |
|
|
|
rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, rc->send_seq, buf, (size_t)n); |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:PAUSE fd=%d sid=%08x len=%zd qcnt=%d — normalizer переполнен, пауза", (int)rc->sock, rc->stream_id, n, qcnt); |
|
|
|
rc->send_seq++; |
|
|
|
rc->pause_buf = u_malloc(n); |
|
|
|
|
|
|
|
if (rc->pause_buf) { memcpy(rc->pause_buf, buf, n); rc->pause_len = n; } |
|
|
|
|
|
|
|
else DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:PAUSE malloc(%zd) failed sid=%08x", n, rc->stream_id); |
|
|
|
|
|
|
|
if (rc->read_id) { uasync_remove_socket_t(rc->ua, rc->sock); rc->read_id = NULL; } |
|
|
|
|
|
|
|
rc->pause_backoff = 50; |
|
|
|
|
|
|
|
rc->pause_timer = uasync_set_timeout(rc->ua, rc->pause_backoff, rc, rp_sock_pause_cb, "rp_pause"); |
|
|
|
|
|
|
|
return; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:RECV fd=%d sid=%08x len=%zd total=%d", (int)rc->sock, rc->stream_id, n, rp_conn_total(rc)); |
|
|
|
|
|
|
|
rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, rc->send_seq, buf, (size_t)n); |
|
|
|
|
|
|
|
rc->send_seq++; |
|
|
|
} else if (n == 0) { |
|
|
|
} else if (n == 0) { |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:EOF fd=%d sid=%08x total=%d cli_closed=%d sock_closed=%d", |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:EOF fd=%d sid=%08x total=%d cli_closed=%d sock_closed=%d", |
|
|
|
(int)rc->sock, rc->stream_id, rp_conn_total(rc), rc->cli_closed, rc->sock_closed); |
|
|
|
(int)rc->sock, rc->stream_id, rp_conn_total(rc), rc->cli_closed, rc->sock_closed); |
|
|
|
@ -113,6 +210,11 @@ static void rp_sock_write_cb(socket_t sock, void* arg) { |
|
|
|
rp_conn_total(rc)); |
|
|
|
rp_conn_total(rc)); |
|
|
|
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; |
|
|
|
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; |
|
|
|
if (inst) rp_send_connected(inst, rc->peer_node_id, rc->stream_id, local_port, TCP_PROXY_CONNECTED_OK); |
|
|
|
if (inst) rp_send_connected(inst, rc->peer_node_id, rc->stream_id, local_port, TCP_PROXY_CONNECTED_OK); |
|
|
|
|
|
|
|
if (rc->pending_buf && rc->pending_len > 0) { |
|
|
|
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:FLUSH fd=%d sid=%08x len=%zu", (int)rc->sock, rc->stream_id, rc->pending_len); |
|
|
|
|
|
|
|
rp_sock_send(rc, rc->pending_buf, rc->pending_len); |
|
|
|
|
|
|
|
u_free(rc->pending_buf); rc->pending_buf = NULL; rc->pending_len = 0; |
|
|
|
|
|
|
|
} |
|
|
|
} else { |
|
|
|
} else { |
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:FAIL fd=%d sid=%08x err=%d %s", |
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:FAIL fd=%d sid=%08x err=%d %s", |
|
|
|
(int)rc->sock, rc->stream_id, err, err ? strerror(err) : "unknown"); |
|
|
|
(int)rc->sock, rc->stream_id, err, err ? strerror(err) : "unknown"); |
|
|
|
@ -146,6 +248,11 @@ void rp_conn_free(struct remote_proxy_conn* rc) { |
|
|
|
if (rc->read_id) { uasync_remove_socket_t(rc->ua, rc->sock); rc->read_id = NULL; } |
|
|
|
if (rc->read_id) { uasync_remove_socket_t(rc->ua, rc->sock); rc->read_id = NULL; } |
|
|
|
socket_close_wrapper(rc->sock); rc->sock = SOCKET_INVALID; |
|
|
|
socket_close_wrapper(rc->sock); rc->sock = SOCKET_INVALID; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
if (rc->pending_buf) { u_free(rc->pending_buf); rc->pending_buf = NULL; } |
|
|
|
|
|
|
|
if (rc->out_timer) { uasync_cancel_timeout(rc->ua, rc->out_timer); rc->out_timer = NULL; } |
|
|
|
|
|
|
|
if (rc->out_buf) { u_free(rc->out_buf); rc->out_buf = NULL; } |
|
|
|
|
|
|
|
if (rc->pause_timer) { uasync_cancel_timeout(rc->ua, rc->pause_timer); rc->pause_timer = NULL; } |
|
|
|
|
|
|
|
if (rc->pause_buf) { u_free(rc->pause_buf); rc->pause_buf = NULL; } |
|
|
|
u_free(rc); |
|
|
|
u_free(rc); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@ -201,8 +308,8 @@ int remote_proxy_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* conn, |
|
|
|
if (!inst) { queue_dgram_free(entry); queue_entry_free(entry); return -1; } |
|
|
|
if (!inst) { queue_dgram_free(entry); queue_entry_free(entry); return -1; } |
|
|
|
struct remote_proxy_ctx* ctx = &inst->remote_proxy; |
|
|
|
struct remote_proxy_ctx* ctx = &inst->remote_proxy; |
|
|
|
struct remote_proxy_conn* rc = remote_proxy_find_conn(ctx, stream_id); |
|
|
|
struct remote_proxy_conn* rc = remote_proxy_find_conn(ctx, stream_id); |
|
|
|
if (!rc || rc->sock == SOCKET_INVALID || rc->connected != 1) { |
|
|
|
if (!rc || rc->sock == SOCKET_INVALID) { |
|
|
|
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "RP handle_data drop: нет активного соединения для sid=%08x, шлём CLOSE", stream_id); |
|
|
|
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "RP handle_data drop: нет соединения для sid=%08x, шлём CLOSE", stream_id); |
|
|
|
if (conn) rp_send_msg(inst, conn->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, stream_id, 0, NULL, 0); |
|
|
|
if (conn) rp_send_msg(inst, conn->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, stream_id, 0, NULL, 0); |
|
|
|
queue_dgram_free(entry); queue_entry_free(entry); return -1; |
|
|
|
queue_dgram_free(entry); queue_entry_free(entry); return -1; |
|
|
|
} |
|
|
|
} |
|
|
|
@ -227,14 +334,18 @@ int remote_proxy_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* conn, |
|
|
|
rc->recv_last_seq = seq; |
|
|
|
rc->recv_last_seq = seq; |
|
|
|
} |
|
|
|
} |
|
|
|
size_t data_len = entry->len - TCP_PROXY_HDR_SIZE; |
|
|
|
size_t data_len = entry->len - TCP_PROXY_HDR_SIZE; |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:SEND fd=%d sid=%08x len=%zu total=%d", |
|
|
|
if (rc->connected) { |
|
|
|
(int)rc->sock, stream_id, data_len, rp_conn_total(rc)); |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:SEND fd=%d sid=%08x len=%zu total=%d", |
|
|
|
if (data_len > 0) { |
|
|
|
(int)rc->sock, stream_id, data_len, rp_conn_total(rc)); |
|
|
|
uint8_t* data = entry->dgram + TCP_PROXY_HDR_SIZE; |
|
|
|
if (data_len > 0) rp_sock_send(rc, entry->dgram + TCP_PROXY_HDR_SIZE, data_len); |
|
|
|
ssize_t n = send(rc->sock, data, data_len, MSG_NOSIGNAL); |
|
|
|
} else { |
|
|
|
if (n < 0 && errno != EAGAIN && errno != EWOULDBLOCK) { |
|
|
|
size_t new_len = rc->pending_len + data_len; |
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: send error %s", strerror(errno)); |
|
|
|
uint8_t* new_buf = u_realloc(rc->pending_buf, new_len); |
|
|
|
} |
|
|
|
if (new_buf) { |
|
|
|
|
|
|
|
if (data_len > 0) memcpy(new_buf + rc->pending_len, entry->dgram + TCP_PROXY_HDR_SIZE, data_len); |
|
|
|
|
|
|
|
rc->pending_buf = new_buf; rc->pending_len = new_len; |
|
|
|
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "RP buffer sid=%08x added=%zu total=%zu", stream_id, data_len, new_len); |
|
|
|
|
|
|
|
} else DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "RP buffer realloc(%zu) failed sid=%08x", new_len, stream_id); |
|
|
|
} |
|
|
|
} |
|
|
|
queue_dgram_free(entry); queue_entry_free(entry); |
|
|
|
queue_dgram_free(entry); queue_entry_free(entry); |
|
|
|
return 0; |
|
|
|
return 0; |
|
|
|
|