summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
authorJesper Jensen <jesper@slashwin.dk>2021-10-05 22:48:37 +0200
committerJesper Jensen <jesper@slashwin.dk>2021-10-05 22:48:37 +0200
commit6e6829ad209550fbac0e363f9421c1fba78915b5 (patch)
treea769fdc053e5cd8dabbebe9048725bdc9b6500cc /src
parente3eb0d2eda0bb8c4decf9bc56b025157e70387c2 (diff)
Add getclient timeout
Diffstat (limited to 'src')
-rw-r--r--src/main.c53
-rw-r--r--src/proto.c231
-rw-r--r--src/routing.c25
-rw-r--r--src/routing.h1
4 files changed, 190 insertions, 120 deletions
diff --git a/src/main.c b/src/main.c
index 5094c8c..3cffa29 100644
--- a/src/main.c
+++ b/src/main.c
@@ -17,12 +17,12 @@ void flush_messages(int sfd, struct message* cursor, const struct message* const
}
int main(int argc, char** argv) {
- struct message outbuff[10] = {0};
+ struct message outbuff[32] = {0};
struct dht dht;
- dht.self = (struct nodeid){.inner={0x0034048f, 0x08000020, 0x00888880, 0x02008460, 0x0ab00521}};
+ dht.self = (struct nodeid){.inner={0xebe9bbf1, 0x3cdba6b3, 0x993e0c87, 0x900d5e25}};
struct message* message_cursor = outbuff;
- proto_begin(&dht, time(NULL), &message_cursor, outbuff+10);
+ proto_begin(&dht, time(NULL), &message_cursor, outbuff+32);
flush_messages(dht.sfd, outbuff, message_cursor);
@@ -34,42 +34,41 @@ int main(int argc, char** argv) {
fflush(stdout);
bool timedout = false;
+ time_t next = 0;
if(!dht.pause){
- time_t next = 0;
struct entry* oldest;
routing_oldest(&oldest);
if(oldest != NULL)
next = oldest->expire;
- for(int i = 0; i < MAX_INFLIGHT; i++) {
- if(!dht.reqalloc[i])
- continue;
+ } else {
+ dbg("DHT timeout is paused");
+ }
- time_t timeout = dht.requestdata[i].timeout;
- if(next == 0 || (timeout != 0 && difftime(timeout, next) < 0))
- next = timeout;
- }
+ for(int i = 0; i < MAX_INFLIGHT; i++) {
+ if(!dht.reqalloc[i])
+ continue;
- if(next != 0) {
- time_t sleepfor = next - time(NULL);
- dbg("Set timeout to %ld", sleepfor);
- struct timeval tv = {
- .tv_sec = sleepfor,
- .tv_usec = 0,
- };
- if(tv.tv_sec <= 0) {
- timedout = true;
- } else {
- setsockopt(dht.sfd, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv));
- }
+ time_t timeout = dht.requestdata[i].timeout;
+ if(next == 0 || (timeout != 0 && difftime(timeout, next) < 0))
+ next = timeout;
+ }
+
+ if(next != 0) {
+ time_t sleepfor = next - time(NULL);
+ dbg("Set timeout to %ld", sleepfor);
+ struct timeval tv = {
+ .tv_sec = sleepfor,
+ .tv_usec = 0,
+ };
+ if(tv.tv_sec <= 0) {
+ timedout = true;
} else {
- fatal("We have to wait some amount of time right now");
+ setsockopt(dht.sfd, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv));
}
-
- } else {
- printf("DHT timeout is paused");
}
+
// Try to receive some data, this is a blocking call
struct sockaddr_storage remote;
socklen_t remote_len = sizeof(remote);
diff --git a/src/proto.c b/src/proto.c
index 520f766..8a6f0e8 100644
--- a/src/proto.c
+++ b/src/proto.c
@@ -132,7 +132,6 @@ int create_ping(char* buff, size_t* buff_len, struct nodeid* self, struct nodeid
}
int send_ping(struct dht* dht, struct nodeid* expected, time_t now, bool node_is_new, const struct sockaddr* dest_addr, socklen_t dest_len, struct msgbuff* msgbuff) {
-
if(*msgbuff->messages >= msgbuff->messages_end)
return PROTO_ENOREQ;
struct message* message = *msgbuff->messages;
@@ -147,8 +146,9 @@ int send_ping(struct dht* dht, struct nodeid* expected, time_t now, bool node_is
message->dest_len = dest_len;
struct ping* data = &dht->requestdata[reqId].cont.ping;
- if(!node_is_new)
+ if(!node_is_new) {
data->remote_id = *expected;
+ }
data->is_new = node_is_new;
data->attempt = 0;
@@ -182,6 +182,8 @@ PROCESS_TIMEOUT(getclient_timeout) {
if(cont->ping.is_new)
return 0;
+ dbg("Discarding node that didn't respond");
+
routing_remove(&cont->ping.remote_id);
return 0;
}
@@ -306,18 +308,17 @@ PROCESS_REPONSE(getclient_response) {
// The response was good, so save the node
if(cont->ping.is_new) {
struct entry* entry;
- if(!routing_offer(&id, &entry)) {
+ if(routing_offer(&id, &entry)) {
+ struct sockaddr_in* ipv4 = (struct sockaddr_in*)remote;
+ entry->addr.ip = ipv4->sin_addr.s_addr;
+ entry->addr.port = ipv4->sin_port;
+ entry->expire = now + PROTO_UNCTM;
+ } else {
dbg("We are no longer interested");
- return 0;
}
-
- struct sockaddr_in* ipv4 = (struct sockaddr_in*)remote;
- entry->addr.ip = ipv4->sin_addr.s_addr;
- entry->addr.port = ipv4->sin_port;
- entry->expire = now + PROTO_UNCTM;
} else {
- fatal("Ping has not been implemented");
- // Touch the routing entry
+ struct entry* entry = routing_get(&id);
+ entry->expire = now + PROTO_UNCTM;
}
// Fan out the search if the results were interesting
@@ -414,6 +415,96 @@ void proto_end(struct dht* dht) {
close(dht->sfd);
}
+int handle_packet(struct dht* dht, time_t now, enum commandType type, char* transaction, size_t transaction_len, char* query, size_t query_len, char* packet, size_t packet_len, struct sockaddr_in* remote, socklen_t remote_len, struct msgbuff* msgbuff) {
+ if(type == CT_RESPONSE) {
+ uint32_t transaction_number;
+
+ if(transaction == NULL) {
+ err("DISCARD: No transaction in response");
+ return 0;
+ }
+
+ // Temporary null terminate the string to parse the number without a copy
+ char* end;
+ dbg("Transaction %s", transaction);
+ transaction_number = strtol(transaction, &end, 10);
+
+ if(end != transaction+transaction_len) {
+ err("DISCARD: Transaction id is not a number %.*s", (int)transaction_len, transaction);
+ return 0;
+ }
+
+ uint16_t reqId;
+ if(!find_req(dht, transaction_number, &reqId)) {
+ err("DISCARD: unknown transaction id %d", transaction_number);
+ return 0;
+ }
+ dbg("Transaction id matches request %d", reqId);
+
+ if(sockaddr_cmp((struct sockaddr*)&dht->requestdata[reqId].addr, (struct sockaddr*)remote) != 0) {
+ err("DISCARD: Unexpected IP for valid transaction");
+ return 0;
+ }
+
+ dht->pause = false;
+ int rc = dht->requestdata[reqId].fun(dht, now, &dht->requestdata[reqId].cont, packet, packet_len, dht->sfd, (struct sockaddr*)remote, remote_len, msgbuff);
+ if(rc == PROTO_ENOREQ) {
+ dht->pause = true;
+ } else if(rc == PROTO_EDISC) {
+ return 0;
+ }
+
+ dht->requestdata[reqId].fun = NULL;
+ dht->requestdata[reqId].timeout_fun = NULL;
+ dht->requestdata[reqId].timeout = 0;
+ dht->reqalloc[reqId] = false;
+ } else if(type == CT_QUERY) { // Must be a query
+ if(query == NULL)
+ fatal("No query function in query request");
+ if(transaction == NULL)
+ fatal("No transaction in request");
+
+ assert(strlen(query) == query_len);
+
+ assert(*msgbuff->messages < msgbuff->messages_end);
+ struct message* message = *msgbuff->messages;
+
+ char* end = message->payload+128;
+ char* cursor = message->payload;
+
+ int rc = snprintf(cursor, end-cursor , "d1:t%ld:", transaction_len);
+ if(rc < 0)
+ fatal("No space for response");
+ cursor += rc;
+ memcpy(cursor, transaction, transaction_len);
+ cursor += transaction_len;
+ rc = snprintf(cursor, end-cursor, "1:y1:r1:r");
+ if(rc < 0)
+ fatal("No space for response");
+ cursor += rc;
+
+ rc = handle_request(&dht->self, query, packet, packet_len, &cursor, end-cursor-1);
+ if(rc == QUERY_EUNK) {
+ dbg("DISCARD: Unknown query method");
+ return 0;
+ } else if(rc != 0) fatal("Error handling request");
+
+ rc = snprintf(cursor, end-cursor, "e");
+ if(rc < 0)
+ fatal("No space for response");
+ cursor += rc;
+ assert(cursor < end);
+
+ message->payload_len = cursor - message->payload;
+
+ memcpy(&message->dest, remote, remote_len);
+ message->dest_len = remote_len;
+ (*msgbuff->messages)++;
+ }
+
+ return 0;
+}
+
int proto_run(struct dht* dht, char* buff, size_t recv_len, struct sockaddr_in* remote, socklen_t remote_len, time_t now, struct message** output, const struct message* const output_end) {
struct msgbuff msgbuff = {
output,
@@ -431,6 +522,10 @@ int proto_run(struct dht* dht, char* buff, size_t recv_len, struct sockaddr_in*
int rc = dht->requestdata[i].timeout_fun(dht, &dht->self, now, &dht->requestdata[i].cont, &msgbuff);
+ if(rc == PROTO_ENOREQ) {
+ dht->pause = true;
+ return 0;
+ }
if(rc != PROTO_EDISC) {
dht->requestdata[i].fun = NULL;
dht->requestdata[i].timeout_fun = NULL;
@@ -446,7 +541,7 @@ int proto_run(struct dht* dht, char* buff, size_t recv_len, struct sockaddr_in*
while(oldest != NULL) {
if(difftime(now, oldest->expire) < 0)
break;
- dbg("============ Ping uncertain node");
+ dbg("Node becomes uncertain");
struct sockaddr_in dest = {0};
dest.sin_family = AF_INET;
@@ -461,7 +556,7 @@ int proto_run(struct dht* dht, char* buff, size_t recv_len, struct sockaddr_in*
return 0;
}
- oldest->expire = now + 30;
+ oldest->expire = 0;
routing_oldest(&oldest);
}
@@ -534,90 +629,40 @@ int proto_run(struct dht* dht, char* buff, size_t recv_len, struct sockaddr_in*
return 0;
}
- if(type == CT_RESPONSE) {
- uint32_t transaction_number;
-
- if(!transaction_set) {
- err("DISCARD: No transaction in response");
- return 0;
- }
-
- // Temporary null terminate the string to parse the number without a copy
- char* end;
- dbg("Transaction %s", transaction);
- transaction_number = strtol(transaction, &end, 10);
-
- if(end != transaction+transaction_len) {
- err("DISCARD: Transaction id is not a number %.*s", (int)transaction_len, transaction);
- return 0;
- }
-
- uint16_t reqId;
- if(!find_req(dht, transaction_number, &reqId)) {
- err("DISCARD: unknown transaction id %d", transaction_number);
- return 0;
- }
- dbg("Transaction id matches request %d", reqId);
+ int rc = handle_packet(dht, now, type, transaction_set ? transaction : NULL, transaction_len, query_set ? query : NULL, query_len, buff, recv_len, remote, remote_len, &msgbuff);
+ assert(rc == 0);
- if(sockaddr_cmp((struct sockaddr*)&dht->requestdata[reqId].addr, (struct sockaddr*)remote) != 0) {
- err("DISCARD: Unexpected IP for valid transaction");
- return 0;
+ {
+ size_t allocated = 0;
+ for(int i = 0; i < MAX_INFLIGHT; i++) {
+ if(dht->reqalloc[i]) {
+ allocated++;
+ }
}
- dht->pause = false;
- int rc = dht->requestdata[reqId].fun(dht, now, &dht->requestdata[reqId].cont, buff, recv_len, dht->sfd, (struct sockaddr*)remote, remote_len, &msgbuff);
- if(rc == PROTO_ENOREQ) {
- dht->pause = true;
- } else if(rc == PROTO_EDISC) {
- return 0;
+ dbg("%d/%d requests pending", allocated, MAX_INFLIGHT);
+ }
+ {
+ int filled;
+ int total;
+#define LFACLEN 32
+ double load_factor[LFACLEN] = {0};
+ routing_status(&filled, &total, load_factor, LFACLEN);
+ dbg("%d/%d nodes in routing table", filled, total);
+
+#define GRAPHY 10
+ for(int y = 0; y < GRAPHY; y++) {
+ for(int x = 0; x < LFACLEN; x++) {
+ if(load_factor[x] > (1.0/GRAPHY) * (GRAPHY-y)) {
+ printf("#");
+ } else {
+ printf(" ");
+ }
+ }
+ printf("|\n");
}
-
- dht->requestdata[reqId].fun = NULL;
- dht->requestdata[reqId].timeout_fun = NULL;
- dht->requestdata[reqId].timeout = 0;
- dht->reqalloc[reqId] = false;
- } else if(type == CT_QUERY) { // Must be a query
- if(!query_set)
- fatal("No query function in query request");
- if(!transaction_set)
- fatal("No transaction in request");
-
- assert(strlen(query) == query_len);
-
- assert(*msgbuff.messages < msgbuff.messages_end);
- struct message* message = *msgbuff.messages;
-
- char* end = message->payload+128;
- char* cursor = message->payload;
-
- int rc = snprintf(cursor, end-cursor , "d1:t%ld:", transaction_len);
- if(rc < 0)
- return EPERM;
- cursor += rc;
- memcpy(cursor, transaction, transaction_len);
- cursor += transaction_len;
- rc = snprintf(cursor, end-cursor, "1:y1:r1:r");
- if(rc < 0)
- return EPERM;
- cursor += rc;
-
- rc = handle_request(&dht->self, query, buff, recv_len, &cursor, end-cursor-1);
- if(rc == QUERY_EUNK) {
- dbg("DISCARD: Unknown query method");
- return 0;
- } else if(rc != 0) fatal("Error handling request");
-
- rc = snprintf(cursor, end-cursor, "e");
- if(rc < 0)
- return EPERM;
- cursor += rc;
- assert(cursor < end);
-
- message->payload_len = cursor - message->payload;
-
- memcpy(&message->dest, remote, remote_len);
- message->dest_len = remote_len;
- (*msgbuff.messages)++;
+#undef GRAPHY
+#undef LFACLEN
}
return 0;
diff --git a/src/routing.c b/src/routing.c
index 1448f55..d77f87c 100644
--- a/src/routing.c
+++ b/src/routing.c
@@ -213,6 +213,9 @@ void routing_oldest(struct entry** dest) {
if(!entry->set)
continue;
+ if(entry->expire == 0)
+ continue;
+
if(*dest == NULL) {
*dest = entry;
continue;
@@ -223,3 +226,25 @@ void routing_oldest(struct entry** dest) {
}
}
}
+
+void routing_status(int* filled, int* size, double* load_factor, size_t load_factor_len) {
+ *size = ROUTINGSIZE;
+
+ *filled = 0;
+ for(size_t i = 0; i < ROUTINGSIZE; i++) {
+ if(table[i].set)
+ (*filled)++;
+ }
+
+ int per_bucket = ROUTINGSIZE / load_factor_len;
+ int overflow = ROUTINGSIZE % load_factor_len;
+ struct entry* table_cursor = table;
+ for(int i = 0; i < load_factor_len; i++) {
+ int is_overflow = i < overflow;
+ for(int j = 0; j < per_bucket + is_overflow; j++) {
+ load_factor[i] += table_cursor->set;
+ table_cursor++;
+ }
+ load_factor[i] /= per_bucket + is_overflow;
+ }
+}
diff --git a/src/routing.h b/src/routing.h
index 3972afa..a7cd8c6 100644
--- a/src/routing.h
+++ b/src/routing.h
@@ -31,6 +31,7 @@ bool routing_interested(struct nodeid* id);
bool routing_offer(struct nodeid* id, struct entry **dest);
void routing_oldest(struct entry** dest);
size_t routing_closest(struct nodeid* needle, size_t n, struct entry** res);
+void routing_status(int* filled, int* size, double* load_factor, size_t load_factor_len);
struct entry* routing_get(struct nodeid* id);
void routing_remove(struct nodeid* self);