Browse Source

remote_proxy: гарантированная доставка CLOSE/ERROR обратно (exit→client) через retry-таймер

- remote_proxy_conn: close_pending, close_timer, close_backoff
- rp_send_close/rp_send_error: при сбое etcp_route_send → close_pending=1 + таймер с backoff
- rp_close_retry_cb: повторная отправка CLOSE/ERROR с экспоненциальным backoff (50→5000 tb)
- rp_conn_free: очистка close_timer
etcp-inflight-fix
Evgeny 4 months ago
parent
commit
e7a9163011
  1. 48
      src/remote_proxy.c
  2. 4
      src/remote_proxy.h

48
src/remote_proxy.c

@ -33,6 +33,7 @@ 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_retry_cb(void* arg);
static void rp_sock_pause_cb(void* arg);
static void rp_close_retry_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);
@ -59,6 +60,44 @@ static int rp_send_connected(struct UTUN_INSTANCE* inst, uint64_t dst, uint32_t
return rp_send_msg(inst, dst, TCP_PROXY_SUBCMD_CONNECTED, sid, buf, 3);
}
static void rp_close_retry_cb(void* arg) {
struct remote_proxy_conn* rc = (struct remote_proxy_conn*)arg;
if (!rc) return;
rc->close_timer = NULL;
if (!rc->close_pending) return;
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
if (!inst) return;
uint8_t subcmd = rc->error ? TCP_PROXY_SUBCMD_ERROR : TCP_PROXY_SUBCMD_CLOSE;
if (rp_send_msg(inst, rc->peer_node_id, subcmd, rc->stream_id, NULL, 0) < 0) {
rc->close_backoff = rc->close_backoff < 5000 ? rc->close_backoff * 2 : 5000;
rc->close_timer = uasync_set_timeout(rc->ua, rc->close_backoff, rc, rp_close_retry_cb, "rp_close_retry");
return;
}
rc->close_pending = 0;
}
static void rp_send_close(struct remote_proxy_conn* rc) {
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
if (!inst) return;
if (rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, rc->stream_id, NULL, 0) < 0) {
rc->close_pending = 1;
rc->close_backoff = 50;
if (!rc->close_timer)
rc->close_timer = uasync_set_timeout(rc->ua, rc->close_backoff, rc, rp_close_retry_cb, "rp_close_retry");
}
}
static void rp_send_error(struct remote_proxy_conn* rc) {
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
if (!inst) return;
if (rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_ERROR, rc->stream_id, NULL, 0) < 0) {
rc->close_pending = 1;
rc->close_backoff = 50;
if (!rc->close_timer)
rc->close_timer = uasync_set_timeout(rc->ua, rc->close_backoff, rc, rp_close_retry_cb, "rp_close_retry");
}
}
static int rp_conn_total(struct remote_proxy_conn* rc) {
if (!rc || !rc->ctx) return 0;
int n = 0; struct remote_proxy_conn* c;
@ -182,9 +221,7 @@ static void rp_sock_read_cb(socket_t sock, void* arg) {
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);
rc->sock_closed = 1;
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
if (inst) rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, rc->stream_id, NULL, 0);
else DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:EOF fd=%d sid=%08x — inst is NULL, can't send CLOSE", (int)rc->sock, rc->stream_id);
rp_send_close(rc);
if (rc->read_id) { uasync_remove_socket_t(rc->ua, rc->sock); rc->read_id = NULL; }
if (rc->cli_closed) rp_conn_free(rc);
} else if (errno != EAGAIN && errno != EWOULDBLOCK && errno != EINTR) {
@ -232,9 +269,7 @@ static void rp_sock_error_cb(socket_t sock, void* arg) {
if (!rc || rc->sock == SOCKET_INVALID) return;
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:ERR fd=%d sid=%08x", (int)rc->sock, rc->stream_id);
rc->error = 1;
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
if (inst) rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_ERROR, rc->stream_id, NULL, 0);
else DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:ERR fd=%d sid=%08x — inst is NULL", (int)rc->sock, rc->stream_id);
rp_send_error(rc);
rp_conn_free(rc);
}
@ -256,6 +291,7 @@ void rp_conn_free(struct remote_proxy_conn* rc) {
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; }
if (rc->close_timer) { uasync_cancel_timeout(rc->ua, rc->close_timer); rc->close_timer = NULL; }
u_free(rc);
}

4
src/remote_proxy.h

@ -39,6 +39,7 @@ struct remote_proxy_conn {
uint8_t sock_closed; // socket EOF received
uint8_t error;
uint8_t connected; // connect() completed successfully
uint8_t close_pending; // CLOSE/ERROR не доставлен, ретрай
uint8_t* pending_buf; // буфер данных пока connect() не завершён
size_t pending_len;
@ -49,6 +50,9 @@ struct remote_proxy_conn {
void* out_timer; // таймер повтора send
int out_backoff; // backoff: 10..5000 tb (1ms..500ms)
void* close_timer; // таймер повтора CLOSE/ERROR
int close_backoff; // backoff: 50..5000 tb (5ms..500ms)
uint8_t* pause_buf; // буфер при паузе чтения (normalizer переполнен)
size_t pause_len;
void* pause_timer; // таймер паузы чтения

Loading…
Cancel
Save