Add a simple rx queue

This commit is contained in:
angt
2016-02-02 21:11:34 +01:00
parent 924df5798f
commit 6561f819f9
2 changed files with 101 additions and 58 deletions

124
mud.c
View File

@@ -34,8 +34,8 @@ struct sock {
struct packet { struct packet {
unsigned char data[MUD_PKT_SIZE]; unsigned char data[MUD_PKT_SIZE];
uint32_t time;
size_t size; size_t size;
// uint32_t time;
// struct path *path; // struct path *path;
}; };
@@ -46,7 +46,8 @@ struct queue {
}; };
struct mud { struct mud {
struct queue queue; struct queue tx;
struct queue rx;
struct sock *sock; struct sock *sock;
struct path *path; struct path *path;
}; };
@@ -281,9 +282,12 @@ struct mud *mud_create (void)
if (!mud) if (!mud)
return NULL; return NULL;
mud->queue.packet = calloc(256, sizeof(struct packet)); mud->tx.packet = calloc(256, sizeof(struct packet));
mud->rx.packet = calloc(256, sizeof(struct packet));
if (!mud->queue.packet) { if (!mud->tx.packet || !mud->rx.packet) {
free(mud->tx.packet);
free(mud->rx.packet);
free(mud); free(mud);
return NULL; return NULL;
} }
@@ -296,11 +300,8 @@ void mud_delete (struct mud *mud)
free(mud); free(mud);
} }
ssize_t mud_recv (struct mud *mud, void *data, size_t size) int mud_pull (struct mud *mud)
{ {
struct sockaddr_storage addr;
socklen_t addrlen = sizeof(addr);
uint32_t now = mud_now(); uint32_t now = mud_now();
if (!now) { if (!now) {
@@ -308,65 +309,100 @@ ssize_t mud_recv (struct mud *mud, void *data, size_t size)
return -1; return -1;
} }
unsigned char buf[2048];
struct sock *sock; struct sock *sock;
ssize_t ret = 0;
for (sock = mud->sock; sock; sock = sock->next) { for (sock = mud->sock; sock; sock = sock->next) {
ret = recvfrom(sock->fd, buf, sizeof(buf), 0, (struct sockaddr *)&addr, &addrlen); unsigned char next = mud->rx.end+1;
if (ret > 0) if (mud->rx.start == next)
break; return 0;
}
struct packet *packet = &mud->rx.packet[mud->rx.end];
struct sockaddr_storage addr;
socklen_t addrlen = sizeof(addr);
ssize_t ret = recvfrom(sock->fd, packet->data, sizeof(packet->data),
0, (struct sockaddr *)&addr, &addrlen);
if (ret<=0) if (ret<=0)
return ret; continue;
if (ret <= 4)
return 0;
struct path *path = mud_new_path(mud, sock->fd, &addr, addrlen); struct path *path = mud_new_path(mud, sock->fd, &addr, addrlen);
if (!path) if (!path)
return -1; return -1;
uint32_t send_now = mud_read32(buf); uint32_t send_now = mud_read32(packet->data);
if (!send_now) { if (!send_now) {
send_now = mud_read32(&buf[4]); send_now = mud_read32(&packet->data[4]);
path->dt = mud_read32(&buf[8]); path->dt = mud_read32(&packet->data[8]);
path->rtt = now-send_now; path->rtt = now-send_now;
errno = EAGAIN; continue;
return -1;
} }
if (path->recv_count == 256) { if (path->recv_count == 256) {
unsigned char reply[3*4]; unsigned char reply[3*4];
uint32_t dt = (now-path->recv_time)>>8; uint32_t dt = (now-path->recv_time)>>8;
path->recv_count = 0; path->recv_count = 0;
path->recv_time = now; path->recv_time = now;
memset(reply, 0, 4); memset(reply, 0, 4);
memcpy(&reply[4], buf, 4); memcpy(&reply[4], packet->data, 4);
mud_write32(&reply[8], dt); mud_write32(&reply[8], dt);
mud_send_path(path, reply, sizeof(reply)); mud_send_path(path, reply, sizeof(reply));
} else { } else {
path->recv_count++; path->recv_count++;
} }
memcpy(data, &buf[4], ret-4); packet->size = ret;
// packet->time = now;
return ret-4; mud->rx.end = next;
} }
void mud_flush (struct mud *mud, uint32_t time) return 0;
}
ssize_t mud_recv (struct mud *mud, void *data, size_t size)
{ {
while (mud->queue.start != mud->queue.end) { if (size+4 < MUD_PKT_SIZE) {
struct packet *packet = &mud->queue.packet[mud->queue.start]; errno = EMSGSIZE;
return -1;
}
if (packet->time > time) if (mud->rx.start == mud->rx.end) {
break; errno = EAGAIN;
return -1;
}
mud->queue.start++; struct packet *packet = &mud->rx.packet[mud->rx.start];
memcpy(data, &packet->data[4], packet->size-4);
mud->rx.start++;
return packet->size-4;
}
int mud_push (struct mud *mud)
{
uint32_t now = mud_now();
if (!now) {
errno = EAGAIN;
return -1;
}
while (mud->tx.start != mud->tx.end) {
struct packet *packet = &mud->tx.packet[mud->tx.start];
// if (packet->time > time)
// break;
mud->tx.start++;
struct path *path = mud->path; struct path *path = mud->path;
ssize_t ret = mud_send_path(path, packet->data, packet->size); ssize_t ret = mud_send_path(path, packet->data, packet->size);
@@ -379,12 +415,14 @@ void mud_flush (struct mud *mud, uint32_t time)
if (path->send_count == 256) { if (path->send_count == 256) {
path->send_count = 0; path->send_count = 0;
path->send_dt = (time-path->send_time)>>8; path->send_dt = (now-path->send_time)>>8;
path->send_time = time; path->send_time = now;
} else { } else {
path->send_count++; path->send_count++;
} }
} }
return 0;
} }
ssize_t mud_send (struct mud *mud, const void *data, size_t size) ssize_t mud_send (struct mud *mud, const void *data, size_t size)
@@ -401,18 +439,20 @@ ssize_t mud_send (struct mud *mud, const void *data, size_t size)
return -1; return -1;
} }
unsigned char next = mud->queue.end+1; unsigned char next = mud->tx.end+1;
if (mud->tx.start == next) {
errno = EAGAIN;
return -1;
}
struct packet *packet = &mud->tx.packet[mud->tx.end];
if (mud->queue.start != next) {
struct packet *packet = &mud->queue.packet[next];
mud_write32(packet->data, now); mud_write32(packet->data, now);
memcpy(&packet->data[4], data, size); memcpy(&packet->data[4], data, size);
packet->size = size+4; packet->size = size+4;
packet->time = now; // packet->time = now;
mud->queue.end = next; mud->tx.end = next;
}
mud_flush(mud, now);
return size; return size;
} }

3
mud.h
View File

@@ -10,5 +10,8 @@ void mud_delete (struct mud *);
int mud_bind (struct mud *, const char *, const char *); int mud_bind (struct mud *, const char *, const char *);
int mud_peer (struct mud *, const char *, const char *); int mud_peer (struct mud *, const char *, const char *);
int mud_pull (struct mud *);
int mud_push (struct mud *);
ssize_t mud_recv (struct mud *, void *, size_t); ssize_t mud_recv (struct mud *, void *, size_t);
ssize_t mud_send (struct mud *, const void *, size_t); ssize_t mud_send (struct mud *, const void *, size_t);