From 188666dc51233051985715b01631e069a7939442 Mon Sep 17 00:00:00 2001 From: Jesper Jensen Date: Wed, 16 Apr 2025 23:34:03 +0200 Subject: Add a first approximation of a lookup algorithm This isn't really "correct" but it's probably going to work. It's a rough implementation of the Kademlia lookup algorithm, with the modification that we don't keep track of anything behind the frontier. We are therefore counting on peers being ok with us potentially sending them the "same" request multiple times. If they respond in a reasonable time, the double request rate _should_ be minimal, but that of course carries the risk that an already overloaded peer would get even more queries. I don't know. It's a first try at something. I need to see it work before I write it off completely. --- src/main.c | 40 +++++------- src/proto.c | 197 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++-- src/proto.h | 20 +++++- test/proto.c | 119 ++++++++++++++++++++++++++++++++++++ 4 files changed, 344 insertions(+), 32 deletions(-) diff --git a/src/main.c b/src/main.c index a18cb78..267f324 100644 --- a/src/main.c +++ b/src/main.c @@ -129,18 +129,6 @@ void flush_messages(int sfd, struct message* cursor, const struct message* const } } -struct lookup { - struct nodeid target; - - struct nodeid closest[8]; - struct addr closest_addr[8]; - bool closest_valid[8]; // @SLOP: This could be a single word - - uint64_t outstanding; - - time_t wake; -}; - #define OUTBOX_SIZE 32 int main(int argc, char** argv) { srand(time(NULL)); @@ -183,20 +171,20 @@ int main(int argc, char** argv) { proto_begin(&dht, time(NULL), &message_cursor, outbuff+32); flush_messages(dht.sfd, outbuff, message_cursor); - /* struct lookup lookup; */ - /* // Init the lookup */ - /* { */ - /* lookup.wake = 0; */ - /* lookup.target = (struct nodeid){.inner={0x19b8a941, 0x38fa0191, 0x1403fac2, 0x581000ab, 0x19583cda}}; */ - - /* struct entry* entry[8]; */ - /* int found = routing_closest(&lookup.target, 8, entry); */ - /* for(size_t i = 0; i < found; i++) { */ - /* lookup.closest[i] = entry[i]->id; */ - /* lookup.closest_addr[i] = entry[i]->addr; */ - /* lookup.closest_valid[i] = true; */ - /* } */ - /* } */ + // Init the lookup + struct lookup *lookup = &dht.lookup; + { + lookup->timeout = 0; + lookup->target = (struct nodeid){.inner={0x19b8a941, 0x38fa0191, 0x1403fac2, 0x581000ab, 0x19583cda}}; + + struct entry* entry[8]; + int found = routing_closest(&lookup->target, sizeof(entry)/sizeof(entry[0]), entry); + assert(found == 8); + for(size_t i = 0; i < found; i++) { + lookup->closest[i] = entry[i]->id; + lookup->closest_addr[i] = entry[i]->addr; + } + } #define RECV_BUFF_SIZE 4096 char buff_storage[RECV_BUFF_SIZE+1]; diff --git a/src/proto.c b/src/proto.c index c1cc392..d5d377c 100644 --- a/src/proto.c +++ b/src/proto.c @@ -157,17 +157,14 @@ bool find_req(struct dht* dht, uint32_t transId, uint16_t* reqId) { return dht->reqalloc[transId]; } -struct msgbuff { - struct message** messages; - const struct message* const messages_end; -}; - #define PROTO_EDISC 1 #define PROTO_ENOREQ 2 PROCESS_REPONSE(getclient_response); PROCESS_TIMEOUT(getclient_timeout); +PROCESS_REPONSE(lookup_response); + // Number of nodeid bits #define IDBITS 160 #if IDBITS > RAND_MAX @@ -205,6 +202,37 @@ int write_find_node(char* buff, size_t* buff_len, struct nodeid* self, struct no return 0; } +int send_lookup(struct dht* dht, struct nodeid* target, time_t now, 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; + + uint16_t reqId; + if(!alloc_req(dht, &reqId)) { + return PROTO_ENOREQ; + } + + memcpy(&message->dest, dest_addr, dest_len); + message->dest_len = dest_len; + + dht->requestdata[reqId].cont.lookup = &dht->lookup; + + dht->requestdata[reqId].fun = &lookup_response; + dht->requestdata[reqId].timeout = now + PROTO_TMOUT; + dht->requestdata[reqId].timeout_fun = NULL; + memcpy(&dht->requestdata[reqId].addr, dest_addr, dest_len); + dht->requestdata[reqId].addr_len = dest_len; + + message->payload_len = sizeof(message->payload); + int rc = write_find_node(message->payload, &message->payload_len, &dht->self, target, reqId); + if(rc != 0) { + return rc; + } + (*msgbuff->messages)++; + + return 0; +} + 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; @@ -245,6 +273,165 @@ int send_ping(struct dht* dht, struct nodeid* expected, time_t now, bool node_is return 0; } +PROCESS_REPONSE(lookup_response) { + struct benc_node stream[256]; + struct bcursor bcursor; + bcur_open(&bcursor, packet, packet+packet_len, stream, 256); + + if(bcursor.end - bcursor.readhead <= 0) { + fatal("Response too short"); + } + + struct nodeid id; + uint8_t nodes_len; + struct nodeid nodes[8]; + struct in_addr ips[8]; + uint16_t ports[8]; + + // Read the payload + { + // Check that we have a dict + if(bcursor.readhead->type != BNT_DICT) { + fatal("Response is not a dict"); + } + bcur_next(&bcursor, 1); + + bcur_find_key(&bcursor, (const enum benc_nodetype[]){BNT_STRING}, (const char*[]){"r"}, (const size_t[]){1}, 1); + // Skip the key + bcur_next(&bcursor, 1); + + if(bcursor.readhead->type != BNT_DICT) { + fatal("Wrong value type for response"); + } + + // Skip the dict element + bcur_next(&bcursor, 1); + + uint8_t parts = 0; + while(bcursor.readhead->type != BNT_END) { + switch(bcur_find_key(&bcursor, (const enum benc_nodetype[]){BNT_STRING, BNT_STRING}, (const char*[]){"nodes", "id"}, (const size_t[]){5, 2}, 2)) { + case 0: + // Skip the key + bcur_next(&bcursor, 1); + + if(bcursor.readhead->type != BNT_STRING) { + fatal("Nodes must be a string"); + } + + if((bcursor.readhead->size % 26) != 0) { + fatal("Nodes string value must be a multiple of 26"); + } + + nodes_len = MIN(bcursor.readhead->size/26, 8); + for(int i = 0; i < nodes_len; i++) { + memcpy(nodes+i, bcursor.readhead->loc+(26*i), 20); + memcpy(ips+i, bcursor.readhead->loc+(26*i)+20, 4); + memcpy(ports+i, bcursor.readhead->loc+(26*i)+24, 2); + } + + parts++; + + // Skip the value + bcur_next(&bcursor, 1); + break; + case 1: + // Skip the key + bcur_next(&bcursor, 1); + + if(bcursor.readhead->type != BNT_STRING) { + fatal("Wrong value type for response"); + } + + if(bcursor.readhead->size != 20) { + fatal("remote node id was not 20 bytes long"); + } + + memcpy(&id, bcursor.readhead->loc, 20); + + parts++; + + // Skip the value + bcur_next(&bcursor, 1); + break; + case -BENC_EBADP: + err("Bad Dictionary, Discard packet"); + return PROTO_EDISC; + } + } + + if(parts < 2) { + err("Response didn't contain nodes and id"); + return PROTO_EDISC; + } + } + + uint8_t my_score = prefix(&id, &cont->lookup->target); + + // We need to find the best match that we are still better than + uint32_t match_i = 0; + uint8_t match_score = 0; + for(size_t i = 0; i < 8; i++) { + dbg("PORT %d", cont->lookup->closest_addr[i].port); + if(cont->lookup->closest_addr[i].port == 0) { + match_i = i; + match_score = UINT8_MAX; // Bogus value to signal that we found something + break; + } + + uint8_t their_score = prefix(&cont->lookup->closest[i], &cont->lookup->target); + if(their_score > match_score && my_score > their_score) { + match_score = their_score; + match_i = i; + } + } + + if(match_score != 0) { + cont->lookup->closest[match_i] = id; + + struct sockaddr_in* ipv4 = (struct sockaddr_in*)remote; + cont->lookup->closest_addr[match_i].ip = ipv4->sin_addr.s_addr; + cont->lookup->closest_addr[match_i].port = ipv4->sin_port; + } else { + dbg("Discarding response from node behind the frontier"); + } + + uint8_t worst_match = UINT8_MAX; + for(size_t i = 0; i < 8; i++) { + if(cont->lookup->closest_addr[i].port == 0) { + worst_match = 0; + break; + } + + worst_match = MIN(worst_match, prefix(&cont->lookup->closest[i], &cont->lookup->target)); + } + + // Fan out the search if the nodes are better than the worst one in the frontier + for(uint8_t i = 0; i < nodes_len; i++) { + // Don't fan out to anything that is a worse match than our current frontier + if(prefix(&nodes[i], &cont->lookup->target) <= worst_match) continue; + + // @ROBUST: Some nodes report a bunch of nodes in the same ip. Maybe we + // could check for that here + + struct sockaddr_in dest = { + .sin_family = AF_INET, + .sin_addr = ips[i], + .sin_port = ports[i], + }; + + int rc = send_lookup(dht, &cont->lookup->target, now, (struct sockaddr*)&dest, sizeof(struct sockaddr_in), msgbuff); + if(rc == PROTO_ENOREQ) { + return rc; + } else if(rc != 0) { + fatal("failed %d", rc); + } + + cont->lookup->timeout = now; + } + + return 0; +} + PROCESS_TIMEOUT(getclient_timeout) { // @HACK: This really sucks. maybe we should just pass in the request id size_t reqId = (typeof(dht->requestdata[0])*)((void*)cont - offsetof(typeof(dht->requestdata[0]), cont)) - dht->requestdata; diff --git a/src/proto.h b/src/proto.h index 0336c0c..0e1bcbc 100644 --- a/src/proto.h +++ b/src/proto.h @@ -21,6 +21,15 @@ struct tokens { size_t head; }; +struct lookup { + struct nodeid target; + + struct nodeid closest[8]; + struct addr closest_addr[8]; + + time_t timeout; +}; + void token_create(struct tokens* tokens, time_t now, struct addr* remote, char* token); int token_validate(struct tokens* tokens, time_t now, struct addr* remote, char* token); @@ -40,10 +49,15 @@ struct ping { union message_cont { struct ping ping; + struct lookup *lookup; }; struct dht; -struct msgbuff; +struct msgbuff { + struct message** messages; + const struct message* const messages_end; +}; + #define PROCESS_REPONSE(NAME) int (NAME)(struct dht* dht, time_t now, union message_cont* cont, char* packet, size_t packet_len, int socket, struct sockaddr* remote, socklen_t remote_len, struct msgbuff* msgbuff) typedef PROCESS_REPONSE(resp); @@ -67,6 +81,8 @@ struct dht { union message_cont cont; } requestdata[MAX_INFLIGHT]; + struct lookup lookup; + time_t wake; struct tokens tokens; }; @@ -78,6 +94,8 @@ struct message { socklen_t dest_len; }; +int send_lookup(struct dht* dht, struct nodeid* target, time_t now, const struct sockaddr* dest_addr, socklen_t dest_len, struct msgbuff* msgbuff); + void proto_begin(struct dht* dht, time_t now, struct message** output, const struct message* const output_end); int proto_run(struct dht* dht, char* buffer, size_t buffer_len, struct sockaddr_in* remote, socklen_t remote_len, time_t now, struct message** output, const struct message* const output_end); void proto_end(struct dht* dht); diff --git a/test/proto.c b/test/proto.c index 54cf158..1c71d6e 100644 --- a/test/proto.c +++ b/test/proto.c @@ -690,3 +690,122 @@ void test_node_closer_to_infohash_is_discovered() { proto_end(&dht); } + +void test_lookup_response() { + struct message outbuff[10] = {0}; + time_t now = 0; + + struct sockaddr_in remote; + + struct dht dht; + dht.self = (struct nodeid){.inner={0x42424242, 0x42424242, 0x42424242, 0x42424242, 0x42424242}}; + routing_init(&dht.self); + + + // There's no requirement that nodes we issue lookups against are in our + // routing table. We therefore don't insert anything in there. + + { + struct message* message_cursor = outbuff; + proto_begin(&dht, 0, &message_cursor, outbuff+2); + // We just ignore the first ping. The bootstrap node never responds. + } + now += 10; + + // There's no fixed api for issuing a lookup command. We just fill in the + // structure and issue the requests + { + struct lookup *lookup = &dht.lookup; + lookup->timeout = now; + lookup->target = (struct nodeid){.inner={0x61616161, 0x61616161, 0x61616161, 0x61616161, 0x61616161}}; + + for(size_t i = 0; i < 8; i++) { + // Setting the port to 0 indicates that this slot is unallocated + lookup->closest_addr[i].port = 0; + } + } + + { + // After setting up the command, we issue a request to a node. + remote = (struct sockaddr_in){ + .sin_family = AF_INET, + .sin_addr = {0xFFFFFFFF}, + .sin_port = 1, + }; + struct message* message_cursor = outbuff; + int rc = send_lookup(&dht, &dht.lookup.target, now, (struct sockaddr*)&remote, sizeof(remote), &(struct msgbuff){&message_cursor, outbuff+2}); + + TEST_ASSERT_EQUAL(rc, 0); + TEST_ASSERT_EQUAL_PTR(message_cursor, outbuff+1); + TEST_ASSERT_EQUAL(91, outbuff[0].payload_len); + TEST_ASSERT_EQUAL_CHAR_ARRAY("d1:ad2:id20:BBBBBBBBBBBBBBBBBBBB6:target20:aaaaaaaaaaaaaaaaaaaae1:q9:find_node1:t1:11:y1:qe", outbuff[0].payload, 91); + TEST_ASSERT_EQUAL(dht.lookup.timeout, now); + } + now += 1; + + { + // The node responds. This should add it to our frontier since it's + // empty. It should also fan out the search into what it returns since + // we still have empty spots after adding this one. + char buff[] = "d1:y1:r1:t1:11:rd2:id20:CBBBBBBBBBBBBBBBBBBB5:nodes26:aBBBBBBBBBBBBBBBBBBB\xFF\xFF\xFF\xFF\x00\x01""ee"; + struct message* message_cursor = outbuff; + int rc = proto_run(&dht, buff, sizeof(buff), (struct sockaddr_in*)&remote, sizeof(remote), now, &message_cursor, outbuff+2); + + TEST_ASSERT_EQUAL(rc, 0); + TEST_ASSERT_EQUAL_CHAR_ARRAY(&dht.lookup.closest[0], "CBBBBBBBBBBBBBBBBBBB", 20); + + TEST_ASSERT_EQUAL_PTR(message_cursor, outbuff+1); + TEST_ASSERT_EQUAL(91, outbuff[0].payload_len); + TEST_ASSERT_EQUAL_CHAR_ARRAY("d1:ad2:id20:BBBBBBBBBBBBBBBBBBBB6:target20:aaaaaaaaaaaaaaaaaaaae1:q9:find_node1:t1:21:y1:qe", outbuff[0].payload, 91); + TEST_ASSERT_EQUAL(dht.lookup.timeout, now); + + memcpy(&remote, &outbuff[0].dest, sizeof(remote)); + } + + // In actual use, we would have a lot more network/public api traffic here + // that would respond to some more pings. I don't want to write that code, + // so instead I'll just fill out some of the internal structured myself. + // + // What we are emulating is that a bunch of nodes responded before the one + // we just asserted above. Those nodes happened to be closer to the final + // target than the outstanding request we have going on. + { + for(size_t i = 0; i < 8; i++) { + dht.lookup.closest[i] = (struct nodeid){.inner={0x42424242, 0x42424242, 0x42424242, 0x42424242, 0x42424242}}; + // The first two bytes match + dht.lookup.closest[i].inner_b[0] = 'a'; + dht.lookup.closest[i].inner_b[1] = 'a'; + dht.lookup.closest[i].inner_b[2] += i; + + dht.lookup.closest_addr[i].ip = 0; + dht.lookup.closest_addr[i].port = 1; + } + } + now += 1; + + { + // The node now finally responds, but woops only the first byte of its + // ID matches. That's worse than the frontier and shouldn't cause any + // addtional adjustment to the frontier. + // It also includes a new node that's also behind the frontier, we + // don't send anything to that one + char buff[] = "d1:y1:r1:t1:21:rd2:id20:aBBBBBBBBBBBBBBBBBBB5:nodes26:aaBBBBBBBBBBBBBBBBBB\xFF\xFF\xFF\xFF\x00\x01""ee"; + struct message* message_cursor = outbuff; + int rc = proto_run(&dht, buff, sizeof(buff), (struct sockaddr_in*)&remote, sizeof(remote), now, &message_cursor, outbuff+2); + + TEST_ASSERT_EQUAL(rc, 0); + // We only need to check this once since we always pick the first slot + // with a given score. It's a little implementation dependant, but it + // beats having 8 asserts. + TEST_ASSERT_EQUAL_CHAR_ARRAY(&dht.lookup.closest[0], "aaBBBBBBBBBBBBBBBBBB", 20); + + // And we didn't fan out to the new node, even though it's actually a better match than any of the ones we have + TEST_ASSERT_EQUAL_PTR(message_cursor, outbuff); + + // Timeout isn't updated since we didn't send anything + // @FRAGILE this has to match the step size of now + TEST_ASSERT_EQUAL(dht.lookup.timeout, now-1); + } + + proto_end(&dht); +} -- cgit v1.2.3