Browse Source

add 16-bit seq to TCP proxy DATA for loss/duplicate detection

congestion
Evgeny 5 months ago
parent
commit
f18037ece2
  1. 32
      src/remote_proxy.c
  2. 9
      src/remote_proxy.h
  3. 29
      src/tcp_proxy.c
  4. 1
      tests/test_remote_proxy.c

32
src/remote_proxy.c

@ -33,7 +33,7 @@ static void rp_sock_error_cb(socket_t sock, void* arg);
static void rp_conn_free(struct remote_proxy_conn* rc); static void rp_conn_free(struct remote_proxy_conn* rc);
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,
uint64_t sid, const uint8_t* data, size_t len) { uint64_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);
if (!e) return -1; if (!e) return -1;
e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + len); e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + len);
@ -41,6 +41,7 @@ static int rp_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd,
e->dgram[0] = ETCP_ID_TCP_PROXY; e->dgram[0] = ETCP_ID_TCP_PROXY;
e->dgram[1] = subcmd; e->dgram[1] = subcmd;
memcpy(e->dgram + 2, &sid, 8); memcpy(e->dgram + 2, &sid, 8);
memcpy(e->dgram + 10, &seq, 2);
if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len); if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len);
e->len = TCP_PROXY_HDR_SIZE + len; e->len = TCP_PROXY_HDR_SIZE + len;
return etcp_route_send(inst, dst, e); return etcp_route_send(inst, dst, e);
@ -50,7 +51,7 @@ static int rp_send_connected(struct UTUN_INSTANCE* inst, uint64_t dst, uint64_t
uint16_t local_port, uint8_t status) { uint16_t local_port, uint8_t status) {
uint8_t buf[3]; uint8_t buf[3];
memcpy(buf, &local_port, 2); buf[2] = status; memcpy(buf, &local_port, 2); buf[2] = status;
return rp_send_msg(inst, dst, TCP_PROXY_SUBCMD_CONNECTED, sid, buf, 3); return rp_send_msg(inst, dst, TCP_PROXY_SUBCMD_CONNECTED, sid, 0, buf, 3);
} }
// ==================================================================== // ====================================================================
@ -62,12 +63,15 @@ static void rp_sock_read_cb(socket_t sock, void* arg) {
uint8_t buf[8192]; ssize_t n = recv(rc->sock, buf, sizeof(buf), 0); uint8_t buf[8192]; ssize_t n = recv(rc->sock, buf, sizeof(buf), 0);
if (n > 0) { if (n > 0) {
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
if (inst) rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, buf, (size_t)n); if (inst) {
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 { } else {
if (n == 0) DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "remote_proxy: EOF stream=%016llx", (unsigned long long)rc->stream_id); if (n == 0) DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "remote_proxy: EOF stream=%016llx", (unsigned long long)rc->stream_id);
else DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: recv error %s", strerror(errno)); else DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: recv error %s", strerror(errno));
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; 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); if (inst) rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, rc->stream_id, 0, NULL, 0);
rc->connected = -1; rc->connected = -1;
rp_conn_free(rc); rp_conn_free(rc);
} }
@ -105,7 +109,7 @@ static void rp_sock_error_cb(socket_t sock, void* arg) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: socket error stream=%016llx", (unsigned long long)(rc ? rc->stream_id : 0)); DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: socket error stream=%016llx", (unsigned long long)(rc ? rc->stream_id : 0));
rc->connected = -1; rc->connected = -1;
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; 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); if (inst) rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, rc->stream_id, 0, NULL, 0);
rp_conn_free(rc); rp_conn_free(rc);
} }
@ -162,6 +166,7 @@ int remote_proxy_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* ent
rc->ctx = ctx; rc->stream_id = stream_id; rc->peer_node_id = src_node_id; rc->ctx = ctx; rc->stream_id = stream_id; rc->peer_node_id = src_node_id;
memcpy(rc->dest_ip, dest_ip, 4); rc->dest_port = dest_port; memcpy(rc->dest_ip, dest_ip, 4); rc->dest_port = dest_port;
rc->ua = inst->ua; rc->sock = SOCKET_INVALID; rc->ua = inst->ua; rc->sock = SOCKET_INVALID;
rc->send_seq = 0; rc->recv_seq_init = 0;
rc->sock = socket(AF_INET, SOCK_STREAM, 0); rc->sock = socket(AF_INET, SOCK_STREAM, 0);
if (rc->sock == SOCKET_INVALID) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: socket() failed"); rp_conn_free(rc); queue_entry_free(entry); queue_dgram_free(entry); return -1; } if (rc->sock == SOCKET_INVALID) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: socket() failed"); rp_conn_free(rc); queue_entry_free(entry); queue_dgram_free(entry); return -1; }
@ -191,6 +196,23 @@ int remote_proxy_handle_data(struct UTUN_INSTANCE* inst, struct ll_entry* entry,
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) { queue_entry_free(entry); queue_dgram_free(entry); return -1; } if (!rc || rc->sock == SOCKET_INVALID || rc->connected != 1) { queue_entry_free(entry); queue_dgram_free(entry); return -1; }
uint16_t seq; memcpy(&seq, entry->dgram + 10, 2);
if (!rc->recv_seq_init) {
rc->recv_last_seq = seq; rc->recv_seq_init = 1;
} else {
int16_t delta = (int16_t)(seq - rc->recv_last_seq);
if (delta <= 0) {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "remote_proxy: duplicate DATA seq=%u stream=%016llx, discard", seq, (unsigned long long)stream_id);
queue_entry_free(entry); queue_dgram_free(entry); return 0;
}
if (delta > 1) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: seq gap seq=%u last=%u stream=%016llx, tearing down", seq, rc->recv_last_seq, (unsigned long long)stream_id);
rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, rc->stream_id, 0, NULL, 0);
rp_conn_free(rc);
queue_entry_free(entry); queue_dgram_free(entry); return -1;
}
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;
uint8_t* data = entry->dgram + TCP_PROXY_HDR_SIZE; uint8_t* data = entry->dgram + TCP_PROXY_HDR_SIZE;
ssize_t n = send(rc->sock, data, data_len, MSG_NOSIGNAL); ssize_t n = send(rc->sock, data, data_len, MSG_NOSIGNAL);

9
src/remote_proxy.h

@ -21,9 +21,9 @@ struct remote_proxy_ctx;
#define TCP_PROXY_CONNECTED_OK 0 #define TCP_PROXY_CONNECTED_OK 0
#define TCP_PROXY_CONNECTED_REFUSED 1 #define TCP_PROXY_CONNECTED_REFUSED 1
#define TCP_PROXY_HDR_SIZE 10 #define TCP_PROXY_HDR_SIZE 12 // svc_id(1)+subcmd(1)+stream_id(8)+seq(2)
#define TCP_PROXY_CONNECT_HDR_SIZE 16 #define TCP_PROXY_CONNECT_HDR_SIZE 18 // HDR_SIZE + dest_ip(4)+dest_port(2)
#define TCP_PROXY_CONNECTED_HDR_SIZE 13 #define TCP_PROXY_CONNECTED_HDR_SIZE 15 // HDR_SIZE + local_port(2)+status(1)
struct remote_proxy_conn { struct remote_proxy_conn {
struct remote_proxy_conn* next; struct remote_proxy_conn* next;
@ -37,6 +37,9 @@ struct remote_proxy_conn {
int connect_called; int connect_called;
uint8_t dest_ip[4]; uint8_t dest_ip[4];
uint16_t dest_port; uint16_t dest_port;
uint16_t send_seq;
uint16_t recv_last_seq;
uint8_t recv_seq_init;
}; };
struct remote_proxy_ctx { struct remote_proxy_ctx {

29
src/tcp_proxy.c

@ -80,6 +80,9 @@ struct etcp_transport {
uint64_t remote_node_id; uint64_t remote_node_id;
uint64_t stream_id; uint64_t stream_id;
int connected; int connected;
uint16_t send_seq;
uint16_t recv_last_seq;
uint8_t recv_seq_init;
}; };
// Forward declarations for lwIP callbacks // Forward declarations for lwIP callbacks
@ -576,8 +579,10 @@ static int etcp_transport_send(struct tcp_proxy_transport* t, const uint8_t* dat
if (!e->dgram) { queue_entry_free(e); return -1; } if (!e->dgram) { queue_entry_free(e); return -1; }
e->dgram[0] = ETCP_ID_TCP_PROXY; e->dgram[1] = TCP_PROXY_SUBCMD_DATA; e->dgram[0] = ETCP_ID_TCP_PROXY; e->dgram[1] = TCP_PROXY_SUBCMD_DATA;
memcpy(e->dgram + 2, &et->stream_id, 8); memcpy(e->dgram + 2, &et->stream_id, 8);
memcpy(e->dgram + 10, &et->send_seq, 2);
if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len); if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len);
e->len = TCP_PROXY_HDR_SIZE + len; e->len = TCP_PROXY_HDR_SIZE + len;
et->send_seq++;
int ret = etcp_route_send(et->inst, et->remote_node_id, e); int ret = etcp_route_send(et->inst, et->remote_node_id, e);
return ret; return ret;
} }
@ -589,7 +594,7 @@ static void etcp_transport_close(struct tcp_proxy_transport* t) {
if (e) { if (e) {
e->dgram = u_malloc(TCP_PROXY_HDR_SIZE); e->dgram = u_malloc(TCP_PROXY_HDR_SIZE);
if (e->dgram) { e->dgram[0] = ETCP_ID_TCP_PROXY; e->dgram[1] = TCP_PROXY_SUBCMD_CLOSE; if (e->dgram) { e->dgram[0] = ETCP_ID_TCP_PROXY; e->dgram[1] = TCP_PROXY_SUBCMD_CLOSE;
memcpy(e->dgram + 2, &et->stream_id, 8); e->len = TCP_PROXY_HDR_SIZE; memcpy(e->dgram + 2, &et->stream_id, 8); memset(e->dgram + 10, 0, 2); e->len = TCP_PROXY_HDR_SIZE;
etcp_route_send(et->inst, et->remote_node_id, e); } else queue_entry_free(e); etcp_route_send(et->inst, et->remote_node_id, e); } else queue_entry_free(e);
} }
et->connected = 0; et->connected = 0;
@ -609,6 +614,7 @@ static struct etcp_transport* etcp_transport_create(struct proxy_conn* pc, struc
et->remote_node_id = remote_node_id; et->remote_node_id = remote_node_id;
et->stream_id = ++pc->proxy->next_stream_id; et->stream_id = ++pc->proxy->next_stream_id;
pc->remote_stream_id = et->stream_id; pc->remote_stream_id = et->stream_id;
et->send_seq = 0; et->recv_seq_init = 0;
uint8_t conn_buf[6]; uint8_t conn_buf[6];
memcpy(conn_buf, pc->dest_ip, 4); memcpy(conn_buf + 4, &pc->dest_port, 2); memcpy(conn_buf, pc->dest_ip, 4); memcpy(conn_buf + 4, &pc->dest_port, 2);
struct ll_entry* e = queue_entry_new(0); struct ll_entry* e = queue_entry_new(0);
@ -617,6 +623,7 @@ static struct etcp_transport* etcp_transport_create(struct proxy_conn* pc, struc
if (!e->dgram) { queue_entry_free(e); u_free(et); return NULL; } if (!e->dgram) { queue_entry_free(e); u_free(et); return NULL; }
e->dgram[0] = ETCP_ID_TCP_PROXY; e->dgram[1] = TCP_PROXY_SUBCMD_CONNECT; e->dgram[0] = ETCP_ID_TCP_PROXY; e->dgram[1] = TCP_PROXY_SUBCMD_CONNECT;
memcpy(e->dgram + 2, &et->stream_id, 8); memcpy(e->dgram + 2, &et->stream_id, 8);
memset(e->dgram + 10, 0, 2);
memcpy(e->dgram + TCP_PROXY_HDR_SIZE, conn_buf, 6); memcpy(e->dgram + TCP_PROXY_HDR_SIZE, conn_buf, 6);
e->len = TCP_PROXY_HDR_SIZE + 6; e->len = TCP_PROXY_HDR_SIZE + 6;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy: remote connect stream=%016llx to %d.%d.%d.%d:%d via node %016llx", DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy: remote connect stream=%016llx to %d.%d.%d.%d:%d via node %016llx",
@ -691,6 +698,26 @@ void tcp_proxy_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
if (subcmd == TCP_PROXY_SUBCMD_DATA) { if (subcmd == TCP_PROXY_SUBCMD_DATA) {
struct proxy_conn* pc = find_pc_by_stream(proxy, stream_id); struct proxy_conn* pc = find_pc_by_stream(proxy, stream_id);
if (pc) { if (pc) {
if (pc->transport && pc->transport->ops == &etcp_transport_ops) {
struct etcp_transport* et = (struct etcp_transport*)pc->transport;
uint16_t seq; memcpy(&seq, entry->dgram + 10, 2);
if (!et->recv_seq_init) {
et->recv_last_seq = seq; et->recv_seq_init = 1;
} else {
int16_t delta = (int16_t)(seq - et->recv_last_seq);
if (delta <= 0) {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy: duplicate DATA seq=%u stream=%016llx, discard", seq, (unsigned long long)stream_id);
queue_entry_free(entry); queue_dgram_free(entry); return;
}
if (delta > 1) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy: seq gap seq=%u last=%u stream=%016llx, tearing down", seq, et->recv_last_seq, (unsigned long long)stream_id);
pc->closing_rem = 1;
if (pc->transport) { pc->transport->ops->destroy(pc->transport); pc->transport = NULL; }
queue_entry_free(entry); queue_dgram_free(entry); return;
}
et->recv_last_seq = seq;
}
}
size_t data_len = entry->len - TCP_PROXY_HDR_SIZE; size_t data_len = entry->len - TCP_PROXY_HDR_SIZE;
if (data_len > 0) { if (data_len > 0) {
struct ll_entry* e = entry_from_data(pc->proxy->entry_pool, entry->dgram + TCP_PROXY_HDR_SIZE, (uint16_t)data_len); struct ll_entry* e = entry_from_data(pc->proxy->entry_pool, entry->dgram + TCP_PROXY_HDR_SIZE, (uint16_t)data_len);

1
tests/test_remote_proxy.c

@ -99,6 +99,7 @@ static void monitor(void* arg) {
e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + 6); e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + 6);
e->dgram[0] = ETCP_ID_TCP_PROXY; e->dgram[1] = TCP_PROXY_SUBCMD_CONNECT; e->dgram[0] = ETCP_ID_TCP_PROXY; e->dgram[1] = TCP_PROXY_SUBCMD_CONNECT;
memcpy(e->dgram + 2, &stream_id, 8); memcpy(e->dgram + 2, &stream_id, 8);
memset(e->dgram + 10, 0, 2);
memcpy(e->dgram + TCP_PROXY_HDR_SIZE, payload, 6); memcpy(e->dgram + TCP_PROXY_HDR_SIZE, payload, 6);
e->len = TCP_PROXY_HDR_SIZE + 6; e->len = TCP_PROXY_HDR_SIZE + 6;
etcp_route_send(inst, inst->node_id, e); etcp_route_send(inst, inst->node_id, e);

Loading…
Cancel
Save