I'm back again after a loong break but this time I have done a lot of things apart from what I said in the last blog (obviously 🙃).
This time I'll be discussing on how I wrote the logic for using TCP in an event loop even though it is a stream oriented protocol. So my idea was simple but it took a long time to come, and this time it was purely my idea — I was sleeping and suddenly I woke up in the middle of the night and my brain was like "Hey, what if we write the size of the packet in it?" and I was like "BROO, right now?"
That did not happen exactly 😂. I simply slept without acknowledging the idea and the next day I woke up and again thought for an hour and then I remembered about this idea.
So what was the idea?
The idea is that say I have my TCP packet wherein I have some fixed size data (we'll do variable size later). Now its really easy in UDP just because UDP is packet oriented which means it preserves packet boundaries meanwhile TCP does not (that's one of the reasons why it is so good!)
So how do we get this amazing UDP property from TCP? The idea is simple, write down the packet size on the packet itself.
Say you have the packet struct defined as:
struct myPacket {
size_t packetSize;
int someData;
char someData2;
// ...some more data
};
Now here, we if do sizeof(struct myPacket) we will get the answer as 16 (yes it's not 13 as I expected earlier because of compiler adding some padding to make 13 a multiple of 4 which makes it easier for the CPU to fetch the pages from memory).
And the value of the packetSize member is 16. So now picture this, you're doing
int got = read(fd, buf, ???);
Here you don't know the value you should put into your buffer since you don't know what packet is coming. But now if every packet sent by the other side (which is also your own code, atleast in my case it is 😂) has this size defined as the first member of the struct you could simply do this:
int got = read(fd, size, sizeof(int));
And once you get that, read the value of buf and figure out the packet size like this:
int size;
int got = read(fd, &size, sizeof(int));
if (size == sizeof(struct myPacket)) {
struct myPacket *mp = malloc(size * sizeof(char));
mp->packetSize = size;
got = read(fd, (char *)mp + sizeof(int), size - sizeof(int));
}
Here I've written more code than needed because I wanted to make it look more "complete" 😅
What's happening here is that first we read an int from the stream. now this int contains the size of the packet (ofc, you'd need ntohs here, but I figured that'd make this more hard to read). And then you simply read the rest of the packet in a struct. Now you don't have to read the remaining data at the start of the struct hence the mp + size instead of mp alone. Now I know it would be much better and safer to use offsetof but for an example this works great.
Anyways this was when everything is great and life is simple. But sadly that ain't true (I wish it was 🫠). So now what if the data hasn't arrived yet? I mean the member someData has arrived but someData2 is still travelling on the network. And when you demand for the entire data read does something clever, it simply blocks you. And that is not a good thing when you're in an event loop since the kernel does not know that you're doing multi threading without telling it and one such blocking call is quite deadly since the entire process goes into the blocked state until the packet is not complete.
Think about that for a second, until the entire packet arrives, no heartbeats, no voting, nothing the node can do. For the node it simply underwent time travel 😵. So we do not want to handle time glitches in a system that is already handling a thousand other things. So the kernel gives us a solution, make the read function non blocking (hurray!). I already wrote this in the previous blog but I'll anyways repeat it. So once the read function is non-blocking we have the control to what we want to do except sleeping.
So here's the thing, the Assigner I have in mind acts as a congestion control for the Gateway and the Worker and hence it looks something like this:
So if suppose the diagram was simpler with the Gateway Data (IN) going to Worker Data (OUT). Now in this case if the Gateway sends more data than what I can store (say, my TCP buffer is full as well as the buffer I allocate using malloc) but the worker is slow and hence the buffer in the workers side sending data eventually fills up and no more data can go in. To visualize it:
So here as you see the data flows in chunks, green indicates the data is there but the buffer is not yet full. Then red indicates the buffer is full and there is not space. So now to implement congestion control I need to stop the gateway when the worker's internal buffer fills up. And since we're using epoll how do we do this?
First thing we do is add the events EPOLLIN and EPOLLOUT to both the sockets along with EPOLLET. Now you'll find the descriptions in the man pages of epoll (7) but its better if I maintain the flow here. So EPOLLET simply means that if the buffer was full and now some data is removed tell me and if the buffer was empty and some data has come in tell me. It won't tell you unless the extremes change, i.e., if new data arrives and there was data previously which you forgot to read, you're basically dead meat 😂, epoll in EPOLLET mode won't tell you if something like that happens. It simply ignores that but if you're in the default mode (i.e., you don't have the flag added) you'll get a totally normal behaviour of epoll continuously waking you up when data is available.
However, in this case the best thing is to use EPOLLET with EPOLLIN and EPOLLOUT. And that is better explained with a visual 😁:
There's a ton of steps but believe me its easy once you understand it (wait for two way congestion control 😁)
Step 1: the gateway sends data and the data is stored in the internal buffer of the socket and the EPOLLIN event is triggered for the gateway socket
Step 2: the assigner reads all the data from the gateway and empties out the internal buffer and puts that data into its own buffer.
Step 3: the assigner writes that data to the internal buffer of the worker socket.
Step 4: more data comes from the gateway, and steps 1-3 repeat.
Step 5: Now the worker socket's internal buffer is full hence send returns -1 with errno set as EAGAIN. So now the assigner will simply remove the EPOLLIN event from the gateway socket. So now even if the gateway sends the data, TCP will simply not allow data to flow because its own buffer is full.
Step 6: The gateway is fast and won't stop so the gateway socket's internal socket gets full. But since the assigner does not read the data, TCP by default stops the data transfer until the data is read.
Step 7: Once the worker reads something, its EPOLLOUT event will trigger because the buffer was full and now it is not hence EPOLLOUT. So now the assigner will first empty its own buffer (not visible in the diagram since it would have introduced a lot more steps) into the workers socket. Now if only some data goes in and the send function again returns an error we don't have an option and we simply wait on the EPOLLOUT event again. But if all the data is sent from the assigner's buffer, the assigner will read more data from the gateway socket. And this newly read data will also go to the worker socket.
Step 8: This newly read data might not again be fully absorbed by the send function hence we don't change any events on the sockets.
Step 9: Now the case when we read from the gateway socket and it returns a value less than the size of the buffer we specified. This means that the gateway has stopped sending the data from the other side and the gateway socket is completely drained. Now what? We simply add the EPOLLIN event to the gateway socket. So now the gateway socket has the EPOLLIN and EPOLLOUT events attached to it.
Step 10: Now keep on feeding the worker until the assigner's buffer isn't empty. Once the assigner's buffer is also empty we simply do nothing since that is the last time the EPOLLOUT event will trigger. From then on the buffer will be at the minimum 1 byte not filled unless the gateway starts sending like a freak again 😂
So that's all the steps the visual is trying to explain. Now you should have a question that why we're using EPOLLIN and EPOLLOUT instead of using only the events which are actually needed. So the gateway never uses EPOLLOUT and the worker never uses EPOLLIN then we should keep the other events for these sockets right?
That's where the bi-directional congestion control comes into play 😁
Bi-directional Congestion Control?
So now you know how we can do congestion control in a single direction but a node will always send and receive so things are never so simple 🫠
First you can try to figure out the bi-directional congestion control mechanism yourself but I'll anyways explain it to you.
So as you see in the diagram below you have two sockets each side:
And you have to make sure congestion control in one direction does not interfere with the data flow of the other direction
Firstly let's name the sockets:
-
GIN- Gateway Data IN -
GOUT- Gateway Data OUT -
WIN- Worker Data IN -
WOUT- Worker Data OUT -
GTW- Gateway To Worker data flow -
WTG- Worker To Gateway data flow
Note: Here
GINandGOUTare a single socket because TCP is a duplex connnection but for a better understand I have shown them separately. And the same goes forWINandWOUT
So now if the Worker has a lot of output to send to the gateway your WIN socket will keep on throwing EPOLLIN events whenever you finish reading it. But then at the same time if the WOUT was full and you've kept only EPOLLOUT then you'll be in trouble because then you won't receive EPOLLIN.
So there is a clever strategy I developed on my own (yeah, this time no gemini brainstorming sessions 😂). You always keep EPOLLOUT on all the OUT sockets. So this means we have EPOLLOUT on WOUT and GOUT. Now you should ask:
Why meet why?
This is because having EPOLLOUT along with EPOLLET gives you a solid advantage to handling bi-directional congestion control. Firstly EPOLLOUT is only ever triggered when the socket goes from totally full to a not full state. So you won't be handling the EPOLLOUT event when some data goes out every time but only when the socket was full and not has at mininum a single byte of free space.
So now that EPOLLOUT is not going anywhere the only event we will ever add/remove is EPOLLIN (EPOLLET is the default, like EPOLLOUT it won't go anywhere 😁). So the advantage of only ever changing EPOLLIN is that if I remove the EPOLLIN from WIN and wait for GOUT's EPOLLOUT to trigger, I can still handle EPOLLIN from GIN or wait for WOUT's EPOLLOUT without having to worry about what the other side is doing.
Isn't that a great advantage? I mean for my case it works pretty well 🫠
Let's code it
So for a bi-directional congestion control mechanism I am not sure if a visual will help or just leave you confused but here it is anyways 🙃
This two way congestion control is working in a similar way to the one way congestion control except now its happening at two places simultaneously 😵
Now its easier to see this in code. So first let's assume some things:
There is a socket named
wand a socket namedgand in here both are TCP sockets soWINandWOUTboth arewitself and similarlyGINandGOUTboth aregitself.Both the sockets are added to a
epollloop and the loop is configured in a way that it runs either of the two functions on any event received and also passes the events downresetEpollEventis a function which takes in the socket and the new events to be added to the socket
So first we write the gateway handler (EPOLLIN for now):
void *handle_gateway(void *args) {
int events = *(int *)args;
if (events & EPOLLIN) {
// data received on g
int size = recv(g, buf, sizeof(buf));
sent_size = send(w, buf, size);
if (sent_size != size) {
// all data was not sent, the buffer is full!
resetEpollEvent(g, EPOLLOUT | EPOLLET);
size_to_send = size - sent_size;
}
}
}
And similarly you can guess the EPOLLIN branch for the worker handler:
void *handle_worker(void *args) {
int events = *(int *)args;
if (events & EPOLLIN) {
// data received on w
int size = recv(w, buf, sizeof(buf));
sent_size = send(g, buf, size);
if (sent_size != size) {
// all data was not sent, the buffer is full!
resetEpollEvent(w, EPOLLOUT | EPOLLET);
size_to_send = size - sent_size;
}
}
}
Now these were for the EPOLLIN event but now let's write the gateway to handle the EPOLLOUT event as well:
void *handle_gateway(void *args) {
int events = *(int *)args;
if (events & EPOLLOUT) {
// now g can be fed some data
int size = send(g, buf, size_to_send);
if (size == size_to_send) {
// all data was sent! take more and send
int sent_size = 0;
size = recv(w, buf, sizeof(buf));
while (size != -1) {
sent_size = send(g, buf, size);
if (sent_size != size) {
// the buffer is full again!
size_to_send = size - max(sent_size, 0);
// here max function is implicitly declared and its purpose is to keep sent_size as a minimum 0 in cases where sent_size might become -1
break;
}
size = recv(w, buf, sizeof(buf));
}
if (size == -1) {
// yay! all data was sent!
resetEpollEvent(w, EPOLLIN | EPOLLOUT | EPOLLET);
}
}
// if the size does not match what we want to send, then just keep waiting again
}
}
I know its a bit longer code than what I wrote for EPOLLIN but I assure you it is worth the dive 😁
And now you should try writing the code for the worker's EPOLLOUT branch yourself (Hey, I'm not lazy 😑)
Apart from the EPOLLIN and EPOLLOUT events you also have 3 events to tell you that the connection is cut EPOLLRDHUP, EPOLLHUP and EPOLLERR these three together can detect almost any kind of connection cuts.
EPOLLRDHUPdetects if a socket is closed from the other side of the network as in using theshutdownfunction which I think no one generally uses 😂. If this event occurs, your socket is writable but not readable and the writability is local, the peer won't be able to read what you wrote if it has closed its reading end.EPOLLHUPalso detects if a socket is closed from the other side but this event detects any errors from the socket and if this event happens then your socket is neither readable nor writable.EPOLLERRdetects some error that happened on the socket. Yeah that's it 😂
Now this was congestion control and tcp packet parsing but both aren't independent, you'll have to handle both together so the parsing happens on the data that comes from the data flow mechanism which also handles congestion inside itself 🫠
But now that you do know the concept this should be much easier. And yeah I struggled to come up with the idea first and then struggled more to implement it but then combining both seemed trivial once I had them independently done.
If I can do it, you can too!
— Some genius (Me 😜)
Just joking! The parsing and congestion control should mostly never overlap since congestion control is for a middleman like an assigner and parsing is for a consumer like a worker or a producer like a gateway (could be the other way round too). A consumer always consumes at its own pace and a producer produces at its own pace and hence the middleman needs congestion control.
In my case the worker (consumer) uses a blocking socket hence even if it is slow or fast does not matter since it will experience time glitches which are fine in my case.
Conclusion
I hope you did not get confused and got to learn something good from this blog but I'm happy to answer any questions you face in the chats.
Disclaimer: Meet can make mistakes, always verify with someone/something
The packet parsing was a bit trivial once I added the size member to the head of the struct. But figuring out the congestion control mechanism demanded me to look into TCP's congestion control (which was of very little help) and in general congestion control and mix them up and create a new mechanism. But then that mechanism was, I'd say, not so "good" hence I threw it in the trash and explored distributed algorithms like paxos, raft and distributed systems like bit torrent and one morning my brain wished me a "Good Morning" with this idea and I was surprised I could think of such a good mechanism. (Yeah I'm proud of my creation, for the first time though 😁)
Currently I'm also new to the distributed realm and you're a better expert than me so if you spot any mistake in my approach or anything please do let me know, I'm more than happy to recreate a new mechanism from scratch and throw this in the trash too 😂 but the new one should be worth it!
So, what's next?
Currently I'm learning some Rust cause I heard it is a harder language than C and I just want to see exactly how hard it is. I'm starting with the fundamentals and mostly using the book The Rust Programming Language which is not exactly a book but a webpage. And I think its pretty awesome since this is the first community I saw making such an amazing and deep "documentation".
And the next time I'm going to discuss on linux containers and how to use OverlayFS on the containers. Also I "might" discuss about GARP and OPFS but all depends on time 😁





Top comments (0)