Rendered at 13:37:47 GMT+0000 (Coordinated Universal Time) with Cloudflare Workers.
Animats 4 days ago [-]
Homa has been around for a while. Here's the 2018 paper.[1]
The core idea: When a message arrives at the sender’s transport module, Homa
divides the message into two parts: an initial unscheduled portion (the first RTTbytes bytes), followed by a scheduled portion.
The sender transmits the unscheduled bytes immediately, using
one or more DATA packets. The scheduled bytes are not transmitted until requested explicitly by the receiver using GRANT packets.
So it sends blind for short requests, then needs a go-ahead from the receiver.
That's reasonable when the main application is a remote procedure call. It's reminiscent of QNX's networking protocol, which is also single packet message request/response but can also handle arbitrarily long messages.
What makes this work today is that per-packet processing overhead in hardware switches is low vs. per-byte overhead. In early software driven switches, per-packet overhead tended to dominate, and sending small packets was very inefficient. In modern hardware switches, where FPGAs are doing the processing, the per-packet overhead is low enough that small packets are not inefficient.
It's amusing that web stuff is so bloated today that any transaction under 1MB is considered "small".
So this is not a suitable protocol for open web use.
The packet sizes should be the same for TCP or Homa; the difference is in sender-based vs. receiver-based congestion control.
Procrastes 4 days ago [-]
Upvote for the QNX mention! I loved using QNX for RT control stuff.
Also, thanks for the clear summary and context.
Veserv 4 days ago [-]
Homa is not a good design. [1]
1. No way to detect whole RPC loss. Since there is no outer connection state, if every packet in the send-side of a RPC is lost then there is no way for a server to detect that it should issue a resend. The RPC is just lost to the ether. This affects small messages, like messages that fit in a single packet, more since there are fewer packets in the send-side.
2. Related to the above, there is no builtin encryption support. So, if you want encryption then you need to layer it either above or below.
3. Benchmarked performance is awful. The 60 kB average message case in [2] Table 4 takes 5(!) hyperthreads to average 20 Gbit/s. That is just 4 Gbit/s per hyperthread. Even a totally naive one-packet per system call network protocol design and implementation should get to ~8 Gbit/s per hyperthread. 30 Gbit/s per hyperthread is easy with just a little focus on performance.
4. Despite all the performance design problems in QUIC (though still faster than Homa) it already solves basically every problem Homa is trying to solve in a much cleaner way. Stream IDs correspond to RPC IDs. Stream Max corresponds to Grants. Multiple streams under single Client allows prioritization.
Except you do not randomly lose entire messages. You can compact small messages into packets. You get more precise RTT time allowing more accurate pacing/congestion calculations. You get builtin encryption. It survives ossified middleboxs. It has multiple ack frames/packets reducing ack overhead.
The only real difference is that Homa uses explicit receiver Resend instead of implicit sender Resend. Except that actually consumes significantly more receiver resources in non-trivial loss scenarios especially due to Resend packets only supporting a single Resend span. It also incurs higher latency and has higher requirements on the entire lossy network path due to all the extra transiting data that you do not want to lose.
And those are just serious problems off the top of my head after reloading the RFC into my head. I can come up with some more if needed.
1. Homa detects and recovers from whole-RPC loss just fine; always has.
2. Homa does not have built-in encryption support (any more than TCP does), but encryption can be (and has been) layered on top pretty easily with DTLS. Isn't that what QUIC does also (layer encryption on top of UDP)?
3. Performance: I think you may be misreading the figure. For W4 or W5, Homa uses about 5 cores to get a total of 40 Gbps (20 Gbps in each direction), which is about 8 Gbps/core. I think you'll be hard-pressed to find a software transport that does much better than this (make sure you are measuring mixed workloads with a variety of message sizes; if you only send large messages then you can do much better (and so does Homa)). Unfortunately, all software transport protocols are resource hogs; in my opinion, transport implementation really should move to the NIC.
4. I would be interested to see measurements of QUIC performance running the same workloads as Homa. My understanding is that QUIC has all of the congestion control problems of TCP, so I'd expect it to perform like TCP under mixed workloads (i.e. 10-100x longer tail latency than Homa for short messages). Among other things, I don't believe QUIC uses switch priority queues, so it does not "solve basically every problem Homa is trying to solve".
Veserv 1 days ago [-]
1. Upon careful re-reading I see:
“Ultimate responsibility for RESENDs lies with the client: if it has not received the entire response message within an expected time then it MUST issue RESEND request(s) for the missing part(s) of the response.”
I retract my statement that it does not describe whole RPC loss in protocol.
However, even as someone familiar with network protocols and reading RFCs that is a very oblique and easy to miss way to describe the process. I would recommend a more direct and explicit treatment of the “request lost” case and make clear the “RESEND Length=-1” mechanism as this seems like the only case where that should be relevant. Likely also a state machine since the state flows are different than the regular data loss mechanism; Tx DATA, Tx Lost, Tx RESEND, Rx RPC_UNKNOWN, Tx DATA or something of the like.
Also as I mentioned in another reply [1] the described mechanism has many undesirable propertys.
2. My statement was more about explicit consideration as to how encryption integrates. Sure, you can just layer it below with IPSec or above with TLS, or any other myriad of ways as you could in literally any protocol for transmitting bytes, but TLS over TCP versus QUIC shows how much more painful it is if it was not in the original design.
3. I did not see that it was a duplex benchmark, but that is a very non-standard way to benchmark, so I can not comment as to comparability to standard one-way benchmarks. If a one-way benchmark is equivalent to the sum of the two-way flows, then I will concede that it is comparable to QUIC implementations.
However, my personal experience with designing and implementing transport protocols results in me finding 8 Gbit/s per core in software to be extremely slow which also applies to every QUIC implementation I have ever seen.
4. Congestion control in QUIC, like in TCP, is not fixed though there are recommended implementations. Furthermore, stream prioritization, something not present in TCP due to the single stream nature, is also not fixed. A implementation that offers prioritizing which streams to transmit data for by SRPT would likely see similar outcomes to Homa. I am not aware of any implementation that does so, but it is a freely configurable element of that protocol. I see nothing that would materially stop a comparable solution from being dropped in.
Mainly, unlike TCP, the design of QUIC has knobs in many of the right places so that you can reasonably configure the desired propertys with minimal design clashing.
Really, the main fundamental distinction in steady state operation is the resend mechanism being largely ack driven TCP and QUIC whereas Homa is largely nack driven.
I am of the opinion that ack driven is much superior technically, which is part of why I think Homa is a bad design, but likely for different reasons and considerations as to how it relates to transport protocol design than the designers of QUIC.
This is also not to say that I think QUIC is the best protocol, QUIC has a number of design and performance weaknesses, but it is overall a reasonable and material improvement over TCP and should be considered the standard to beat.
My knowledge of QUIC is pretty superficial, but it sounds like you are saying that congestion control is not a fundamental property of QUIC, so someone could potentially implement Homa's congestion control mechanism in QUIC. If someone were to do that, then of course I'd expect that to perform similarly to Homa.
However, this would create a conflict within QUIC. My understanding is that one of QUIC's primary use cases is WAN connections between browsers and Web servers. In this world Homa's congestion control mechanism does not work; it only works in datacenter environments with speed-of-light delays at most a few 10's of microseconds. Thus, QUIC would need to support multiple congestion control mechanisms, one for WAN and one for datacenter and choose between them as appropriate.
In addition, my understanding is that QUIC is stream-based, which would mean, as with TCP, the protocol doesn't know about message lengths. This makes it hard to prioritize short messages.
And, by the way, as much as you may prefer acks to grants, the grant approach is pretty fundamental to Homa's congestion control.
To me, QUIC and Homa were designed for very different environments. Homa will never compete with QUIC in the WAN, and I think it will be hard for QUIC to compete with Homa in the datacenter (Homa can specialize itself entirely for datacenter apps, whereas QUIC can't).
Veserv 13 hours ago [-]
You should do a deeper dive into QUIC, it has a number of interesting ideas.
1. Having multiple congestion control algorithms is already standard in both QUIC and TCP. As deploying Homa requires you to control both ends or at least understand both ends, in a comparable QUIC deployment you would know they are in a datacenter and can just link the corresponding implementation. Or you could use the configuration parameter handshake mechanism to switch at runtime for certain flows.
2. QUIC is stream-based, but they support fixed/final-sized streams which are logically equivalent to messages. You just generate a disposable fixed-size stream for each “message”. So, it is trivial for the Tx side to know the “message” size and thus do SRPT prioritization. The thing that would be more annoying to support is the Rx side knowing the stream size as size finalization is detected by the Rx side via the FIN bit when transmitted the terminal segment in a stream. You could transmit the terminal segment first, but that is awkward and Homa is superior in that respect. However, Rx size finalization would only really be important for avoiding over-allocation of Rx buffering in QUIC since they do not rely on MAX_STREAM_DATA for flow control whereas the equivalent concept in Homa, GRANT, is used for flow control and thus precise Tx sizing awareness by the Rx side is more important.
3. As mentioned above, Homa GRANT is very similar to QUIC MAX_STREAM_DATA which are basically just extensions of TCP window size to more than one data flow. The main difference is how they are used.
TCP and Homa both use it for Rx resource availability and bandwidth limiting. Homa uses it additionally for data flow prioritization as it has more than one data flow. QUIC uses it for Rx resource availability and implementations can use it for data flow prioritization. QUIC uses MAX_DATA or just plain congestion control bandwidth estimation for bandwidth limiting.
The fundamental usage is Rx resource availability. If you have no Rx buffers available to take the data, then no amount of data sent is meaningful. That is a hard cap. However, using that as a bandwidth limiting mechanism like TCP and Homa requires continuous tight feedback to avoid the Tx from determining there are no Rx resources available while not over-provisioning allowed bandwidth.
To illustrate, suppose I have a 10 GB/s NIC and generate a 1 GB request per second, but my end-to-end is a 1 GB/s link with 1 ms one-way latency. So, this link only handles 1 MB in-flight without buffering. For the sake of simplicity say we pack only 1 KB per packet, so the link can only have 1,000 packets in flight without buffering. How should we manage the GRANT flow?
If we do a initial GRANT of 10 MB, then the Tx will shovel 10 MB into its 10 GB/s NIC which will spit out into the network in 1 ms. This will clog up at the 1 GB/s bottleneck requiring ~10 MB of buffering at the bottleneck. The Tx will then wait for more GRANT as it just bursted it all. After 1 ms the Rx will get the first byte, what should it GRANT and when? If we wait 1 ms, we have received our first MB, should we GRANT 1 MB? Our steady state then looks like bursts of 10 GB/s with 10 MB of buffering at the bottleneck until the request completes.
If we do a initial GRANT of 1 MB then GRANT 1 MB when we get it all, then we stall.
If we do a initial GRANT of 2 MB, then 1 MB every ms then we get 2 MB of buffering at the bottleneck until the request completes.
What is the optimal flow? The optimal outcome is a continuous 1 GB/s flow with one 1 KB packet entering every 1 us. This results in minimal buffering as a result of the transport protocol. For the GRANT mechanism to do this would require a new 1 KB GRANT to arrive every 1 us.
Absent a separate point-to-point bandwidth control system to pace the DATA flow you are not getting that. But once you have that, the GRANT is no longer the bandwidth limiter, it becomes a Rx resource availability and flow prioritization mechanism, just like QUIC. It no longer becomes meaningfully distinct.
KerrAvon 4 days ago [-]
QUIC is much more well known than Homa. Are there downsides for using it for this application?
Veserv 4 days ago [-]
Compared to TCP or Homa at a protocol design level? Not really.
However, software QUIC implementations are generally much slower than software TCP implementations. This is not due to any fundamental protocol design limitations, it is just that most/every QUIC implementation is poorly implemented for maximizing performance.
You can, of course, do multiple times more throughput than QUIC or TCP with a protocol better designed for performance, but that is likely orthogonal to your question.
HackerThemAll 3 days ago [-]
> most/every QUIC implementation is poorly implemented for maximizing performance.
bold claim. Go tell that to Cloudflare, for example. And/or show us your implementation.
Veserv 3 days ago [-]
Third party benchmarking [1] shows MsQuic solidly outperforming quiche by Cloudflare.
The official MsQuic benchmarks [2] tuned for the standard QUIC benchmark [3] with massively over-provisioned hardware does not even reach 8 Gbit/s on large sends with negligible RTT and loss.
If the highest performing public implementation can not even reach a paltry 8 Gbit/s it is fair to say they are all poorly performing.
As for what I would consider “okay” that would be around 30 Gbit/s per core, bottlenecking on decryption (assuming you do not use the allowed null encryption for benchmarking), or bottlenecking on ~1/3 of the bulk 1 Kbyte memcpy() throughput on your hardware, whatever comes first. “Good” is on the order of 100 Gbit/s per core or the same external limiting factors.
1. You have to use timeouts just like TCP has a SYN timeout.
2. TCP doesn't have encryption either. They should probably show DTLS working though.
4. You can't just say QUIC is faster without testing it.
Veserv 3 days ago [-]
1. I do not see where they discuss that in the RFC. While that would work, it sounds awful.
That is a sender-driven resend and as such relies on the sender identifying the “send lost” condition. However, the entire protocol is designed around not doing that and thus you can only feasibly rely on “should have received a reply by now” as your timeout.
But Homa is intended to be a RPC protocol. So you send to server, wait for server to process command, then wait for server to send reply. Your timeout depends on the variable and heterogeneous command processing time.
Even if you were able to give a separate correctly tuned timeout for every possible RPC that is still awful. Any RPC with long processing time should not trigger the timeout until the expected reply time, but if that is far larger than the RTT then you are waiting a tremendous amount of time.
For instance, a select query that is only a few bytes long (and thus fits in one packet) could take seconds on a large database even though the database server is physically nearby and only microseconds away. In that case you would have to wait seconds before timing out instead of just microseconds like a ack-based design could achieve. The worst-case transport latency becomes application-specific instead of related to the application-independent transport parameters like RTT. That is troubling.
2. TCP can be forgiven for that given it predates asymmetric cryptography and even DES. Not considering it relevant in a new protocol this side of the millennium is much less reasonable.
I'm not sure what you're getting at with your point 1. Any implementation of RPCs, over any protocol, faces the same timeout issues. If it takes a long time to service an RPC, the client will naturally wonder if something broke. Either the client has to send some sort of "ping" to make sure the server got the request and is still alive, or the server has to send "keep alive" packets to let the client know it's still alive. If most RPCs finish quickly then the ping interval should be short in order to detect packet drops quickly; but then a large number of pings will be required for RPCs that take a long time to finish.
I don't see anything particularly Homa-specific in these issues. A protocol like TCP that doesn't implement RPCs doesn't have to face these issues; it punts them up to the application, so the application has to deal with them. Homa implements all of this in the transport, so applications don't have to worry about it.
Veserv 14 hours ago [-]
Example: 1 ms RTT with 1000 ms command processing time.
In a lossless RPC the client should expect the reply in 1002 ms. When should the client resend the request if it got lost?
Given that the client would not expect to see the first byte of the reply until 1002 ms and receives no other feedback until then, they should not issue a resend of their request until 1002 ms.
That means if the network lost their first request, they will not send a resend until 1002 ms and then get the reply to their resend at 2004 ms.
In contrast, in a ack based approach, the client would expect a server ack in ~2 ms as the server gives feedback on receipt rather than on processing completion/reply start.
That means if the network lost their first request they will not send a resend until 2 ms and then get the reply to their resend at 1004 ms, nearly half the latency.
In ack based approaches, packet loss requires one additional RTT latency to recover per repeated packet loss. But Homa incurs one additional RTT + processing latency to recover per repeated packet loss.
As a secondary downside, this requires the client to hold Tx resources for RTT + processing time since the client does not know when the server has received all of the data. In ack-based approaches, you learn the server has received all of the request when it acks all of the data in the request which is on the order of the RTT and thus can release Tx resources at that time. This is similar to the reason Homa has ACK/NEED_ACK packets except applied to client side resources instead of just server side resources.
The case you are describing is more like what if the server explodes, but you still need the results/modifications related to your RPC to occur. In that case, whether the Tx got through or not is irrelevant because your RPC still did not complete and you need it to complete, so there is no advantage to the ack-based approach. But by that logic, what if it explodes before persisting changes? Or what if all your servers explode? At some point we should call that out of scope for a transport protocol.
I prefer setting that line at “bytes got from one end to the other”. If you want higher guarantees then you can layer that. As you point out, it is easy to just layer a “server dead” timeout ping which basically solves the problem in the same way Homa needs to.
wmf 3 days ago [-]
You keep talking about Gbits but Homa is about latency.
Veserv 3 days ago [-]
There is nothing in the design that supports low latency relative to a modern protocol like QUIC. Basically the entire “latency” advantage it has over TCP is just that it does not head-of-line block which QUIC also does not do.
And then in basically every other respect it is worse including compatibility.
boesboes 3 days ago [-]
You keep changing the goal posts, typical
tcdent 4 days ago [-]
Why would you run encryption between nodes in an AI cluster?
Because the NSA and similar organizations will 100% infiltrate the physical infrastructure. It’s easier and harder if not impossible to detect.
MrDrMcCoy 3 days ago [-]
So that any sensitive data stupidly uploaded doesn't get leaked, which includes "zero-retention" corporate stuff. Anything that could get sniffed by employees and other agents. It really should be standard practice for any data passing between nodes, even if it's all within a site.
Veserv 3 days ago [-]
Homa was proposed as a datacenter-centric protocol for intra-DC traffic. In that context, with heterogeneous compute loads with various levels of criticality and confidentiality requirements, not considering how to layer encryption better is a design flaw.
locknitpicker 3 days ago [-]
> Why would you run encryption between nodes in an (...) cluster?
This is self-evident, isn't it?
Other than the buzzword factor, why do you think name-dropping AI makes a difference?
preisschild 3 days ago [-]
Because those LLMs often operate on confidential data that should not be leaked?
giovannibonetti 4 days ago [-]
I remember listening to Jane Street’s Ron Minsky on their podcast talking about this a few months ago, how TCP becomes the bottleneck in AI clusters. As an electrical engineer, I remember that circuit switching gave away to packet switching due to very sparse usage of the network when there are many actors going through it. It is not very efficient, but that wide variety of traffic makes it hard to optimize it since the flow patterns are too dynamic. A good analogy with car traffic is that downtown there are so many cars going to a large variety of places, that traffic lights – as inefficient as they are – are a solution that at least works good enough.
On the other hand, if the traffic follows a very predictable pattern, a custom implementation can be much more efficient. Specially nowadays machine learning can find much better solutions through reinforcement learning. And AI cluster data flow is much more predictable than what goes over the internet as a whole.
cma 3 days ago [-]
Google TPUs use optically routable fibers like circuit switching to be able to change the whole topology of the setup depending on training load/matrix dimensions.
smj-edison 3 days ago [-]
I remember reading about their setup and it blew my mind. They have MEMS mirrors that can rotate to redirect one of thousands of fiber lines into any other fiber line by using other optics. That way you can take in like 1000 lines and optically connect it to any of the other 1000 lines.
scotty79 3 days ago [-]
I wonder how did they pull that off. When I imagine end of a fiber optic, the light doesn't come out of it nicely in one direction.
cma 3 days ago [-]
TPU v4: An Optically Reconfigurable Supercomputer for Machine Learning with Hardware Support for Embeddings
Basically they use lens array (one lens per input fiber) to focus each specific beam on its assigned mirror and mirror directs the light to the chosen output that also each has a lens in front of the fiber going out.
scrubs 3 days ago [-]
A few tidbits which may be of interest re homa:
1 - https://news.ycombinator.com/item?id=48287718 is my post on tla formal modeling complete with examples and decent pdf to go with it. Self contained assumes no background on tla
2 - part of the reason for doing (1) was to model homa's basic protocol. See pdf. I did not go into homa details there because I wanted to stay focused on tla.
3 - I have a kernel bypass homa transport api in development
4 - as wonderful as homa is for intra-data-center IO, the srpt part of the protocol is a bit tough to performantly code. More or less homa wants each NIC resource to be managed by one owner who chooses from among all tx streams that could be sent, which one has smallest srpt (shortest remaining processing time). This applies despite the fact client/servers sharing the NIC run in their own threads. Ive been using mpsc memory rings moving the smallest amount of data i can.
throw0101c 4 days ago [-]
With regards to (~6:01) "congestion control is the responsibility of the sender" and "somehow we have to get the sending nodes to stop sending so fast". Does that not exist in Ethernet/RoCE?
> Link Level Flow Control: InfiniBand uses a credit-based algorithm to guarantee lossless HCA-to-HCA communication. RoCE runs on top of Ethernet. Implementations may require lossless Ethernet network for reaching to performance characteristics similar to InfiniBand. Lossless Ethernet is typically configured via Ethernet flow control or priority flow control (PFC). Configuring a Data center bridging (DCB) Ethernet network can be more complex than configuring an InfiniBand network.[19]
> A sending station (computer or network switch) may be transmitting data faster than the other end of the link can accept it. Using flow control, the receiving station can signal the sender requesting suspension of transmissions until the receiver catches up. Flow control on Ethernet can be implemented at the data link layer.
Ethernet flow control doesn't really fix congestion except in very special cases.
Consider a very simple topology:
A C
\ /
S1===S2
/ \
B D
Say hosts A and B are both sending data to C, as fast as they can, via switches S1 and S2 (which are connected via a high-speed link). And say the sum of these two flows is more than the capacity of the link to C.
S2 is receiving packets destined for C faster than it can forward them, but sending an Ethernet pause frame from S2 to S1 is not a very productive way to alleviate the situation, because it also disrupts any traffic that would be bound for D. It just moves the bottleneck elsewhere and causes collateral damage.
stingraycharles 4 days ago [-]
Maybe, but doesn’t RoCE / RDMA handle this at a higher level as well? It’s why PFC and ECN are required for these deployments, which effectively ensure everyone behaves nicely.
I agree that Ethernet flow control is insufficient, but given that NVidia in an act of brilliant foresight acquired Mellanox a decade ago, I’m fairly certain that this is how all these AI clusters are actually deployed, not using TCP, and maybe not even using Ethernet but infiniband instead.
etc-hosts 3 days ago [-]
if you spend the zillion dollars on switches and racks and power and cables and redundant power in your datacenter for infiniband(the only vendor: Nvidia), you get IB.
if not, you get RoCE over ethernet.
throw0101a 3 days ago [-]
> if you spend the zillion dollars on switches and racks and power and cables and redundant power in your datacenter for infiniband(the only vendor: Nvidia), you get IB.
Worth noting that Nvidia/Mellanox is currently the only vendor for IB, but that has not always been the case. I had Intel IB switches for the backend network of my Isilon at my last job (we starting using them when they were 'just' Isilon, then EMC Isilon, and now Dell-EMC Isilon); only post-Dell did they start going Ethernet-only, and now with AI/ML they're re-offering IB as an option (LOL).
So Nvidia is kind of being 'rewarded' for sticking with a technology that everyone gave up on (I guess HPC wasn't big enough to bother with).
Besides IB, the other options being used in Top 500 are Slingshot (Eth-base?) and Omni-path (Category: Interconnect Family):
I hope whoever led the Mellanox acquisition at Nvidia is treated internally like a prophet.
spwa4 3 days ago [-]
... which brings the question, "how is X solved?".
For IB the answer is generally a centralized controller, a lossless flow control algorithm and the credit system with subnet coordination. Ie. the packets that would be dropped due to buffer problems never get sent because the originating node knows it's out of credits. This means packet drops should be nonexistent and so the protocols make the assumption only one packet in a zillion ever gets dropped. On the down side: the tiniest of problems on IB networks have incredible impact (but it's realistic to just not have those)
IB doesn't interface with others. Yes, there's an ethernet bridge but it's not going to work nearly as well once you insert that bit of kit.
In ethernet packets are dropped when congestion is experienced, and you configure through priority queues whose packets get dropped first (you configure priority), and "stuff will sort itself out" like on the internet. RoCE tries to construct IB semantics from this, but there's limits. On the plus side: this is how the internet works and so the whole actually works somewhat reasonably over internet or other long-distance links (and tunnels into Neoclouds and/or Clouds). On the downside it's never going to match IB in speed (it's still going to get to 90% or so link speed, 98% if you just naively check the link utilization and ignore the effect retransmissions will have)
486sx33 4 days ago [-]
[dead]
lokar 4 days ago [-]
The issue is never really with the destination, it’s some link/switch along the way.
When a link becomes saturated working out how to manage that is a hard problem.
I have only used RoCE once at scale, it was really finicky. We would get big waves to pause frames that stalled everything.
wmf 4 days ago [-]
Flow control is better than nothing but it can cause congestion spreading and bufferbloat. QCN, Falcon, and Ultra Ethernet provide much better congestion control for RoCE but they also require newer hardware compared to Homa.
adastra22 4 days ago [-]
Please don’t make the primary link a video.
netsharc 3 days ago [-]
My comment is offtopic, but wow what a cheap intro animation. The lens flare is so fake, you can smell the sloppiness.
One underlying claim here is that latency within the data center will start to matter more and more for future AI compute.
That is relevant for musks in-space computation. There, latency between compute nodes, due to being hundreds of kilometres apart, cannot be low.
This pretty much limits the size of AI models which can be trained on a space datacenter to a single satellite to avoid incurring huge latency overhead both in training and inference.
With the smartest models getting bigger and bigger, that really doesn't look good for space based compute.
rwmj 3 days ago [-]
Putting data centers in space is such obvious nonsense that you don't need to give it serious consideration.
jabl 3 days ago [-]
For another approach that has actually been deployed at scale by OpenAI there's MRC:
MRC is an extension of RoCE that does out-of-order packet delivery, packet spraying across all available paths using SRv6 routing, combination of ECN and packet trimming for congestion control (so sender-based CC like TCP and unlike Homa).
throw0101c 3 days ago [-]
One of the Packet Pushers podcasts just had an interview on this:
> Today's show puts the “heavy” in Heavy Networking with a discussion about Multipath Reliable Connection (MRC). MRC is about getting an even distribution of RDMA traffic across equal cost multipath links in massive data centers, and doing so over plain old lossy Ethernet. Ethan, Drew, and guests discuss how package spraying, resilience to link and fabric failures, and congestion control allow MRC to be a specialized transport for AI workloads.
> Our guests today are from Broadcom: Rip Sohan, Distinguished Engineer; Eric Davis, Master Engineer; and Eric Spada, Technical Director and Distinguished Engineer.
I wonder if, as LLMs get accustomed to using homa or some other optimized protocol for, as the article lists, "chores such as weight gradients, model weights, KV cache entries, and checkpoints", whether we'll start to see TCP as a bottleneck for their post-trained interactions as well, for many of the reasons.
Actually AI is a place where TCP works mostly fine, due to its bandwidth limited nature rather latency limited.
ContinuityLab 3 days ago [-]
A compelling architectural look at moving beyond traditional TCP to optimize latency and throughput for modern distributed AI workloads. Low-level networking is becoming the ultimate bottleneck.
The core idea: When a message arrives at the sender’s transport module, Homa divides the message into two parts: an initial unscheduled portion (the first RTTbytes bytes), followed by a scheduled portion. The sender transmits the unscheduled bytes immediately, using one or more DATA packets. The scheduled bytes are not transmitted until requested explicitly by the receiver using GRANT packets.
So it sends blind for short requests, then needs a go-ahead from the receiver. That's reasonable when the main application is a remote procedure call. It's reminiscent of QNX's networking protocol, which is also single packet message request/response but can also handle arbitrarily long messages.
What makes this work today is that per-packet processing overhead in hardware switches is low vs. per-byte overhead. In early software driven switches, per-packet overhead tended to dominate, and sending small packets was very inefficient. In modern hardware switches, where FPGAs are doing the processing, the per-packet overhead is low enough that small packets are not inefficient.
It's amusing that web stuff is so bloated today that any transaction under 1MB is considered "small". So this is not a suitable protocol for open web use.
[1] https://people.csail.mit.edu/alizadeh/papers/homa-sigcomm18....
Also, thanks for the clear summary and context.
1. No way to detect whole RPC loss. Since there is no outer connection state, if every packet in the send-side of a RPC is lost then there is no way for a server to detect that it should issue a resend. The RPC is just lost to the ether. This affects small messages, like messages that fit in a single packet, more since there are fewer packets in the send-side.
2. Related to the above, there is no builtin encryption support. So, if you want encryption then you need to layer it either above or below.
3. Benchmarked performance is awful. The 60 kB average message case in [2] Table 4 takes 5(!) hyperthreads to average 20 Gbit/s. That is just 4 Gbit/s per hyperthread. Even a totally naive one-packet per system call network protocol design and implementation should get to ~8 Gbit/s per hyperthread. 30 Gbit/s per hyperthread is easy with just a little focus on performance.
4. Despite all the performance design problems in QUIC (though still faster than Homa) it already solves basically every problem Homa is trying to solve in a much cleaner way. Stream IDs correspond to RPC IDs. Stream Max corresponds to Grants. Multiple streams under single Client allows prioritization.
Except you do not randomly lose entire messages. You can compact small messages into packets. You get more precise RTT time allowing more accurate pacing/congestion calculations. You get builtin encryption. It survives ossified middleboxs. It has multiple ack frames/packets reducing ack overhead.
The only real difference is that Homa uses explicit receiver Resend instead of implicit sender Resend. Except that actually consumes significantly more receiver resources in non-trivial loss scenarios especially due to Resend packets only supporting a single Resend span. It also incurs higher latency and has higher requirements on the entire lossy network path due to all the extra transiting data that you do not want to lose.
And those are just serious problems off the top of my head after reloading the RFC into my head. I can come up with some more if needed.
[1] https://github.com/johnousterhout/homa-rfc/blob/main/draft-o...
[2] https://www.usenix.org/system/files/atc21-ousterhout.pdf
1. Homa detects and recovers from whole-RPC loss just fine; always has.
2. Homa does not have built-in encryption support (any more than TCP does), but encryption can be (and has been) layered on top pretty easily with DTLS. Isn't that what QUIC does also (layer encryption on top of UDP)?
3. Performance: I think you may be misreading the figure. For W4 or W5, Homa uses about 5 cores to get a total of 40 Gbps (20 Gbps in each direction), which is about 8 Gbps/core. I think you'll be hard-pressed to find a software transport that does much better than this (make sure you are measuring mixed workloads with a variety of message sizes; if you only send large messages then you can do much better (and so does Homa)). Unfortunately, all software transport protocols are resource hogs; in my opinion, transport implementation really should move to the NIC.
4. I would be interested to see measurements of QUIC performance running the same workloads as Homa. My understanding is that QUIC has all of the congestion control problems of TCP, so I'd expect it to perform like TCP under mixed workloads (i.e. 10-100x longer tail latency than Homa for short messages). Among other things, I don't believe QUIC uses switch priority queues, so it does not "solve basically every problem Homa is trying to solve".
“Ultimate responsibility for RESENDs lies with the client: if it has not received the entire response message within an expected time then it MUST issue RESEND request(s) for the missing part(s) of the response.”
I retract my statement that it does not describe whole RPC loss in protocol.
However, even as someone familiar with network protocols and reading RFCs that is a very oblique and easy to miss way to describe the process. I would recommend a more direct and explicit treatment of the “request lost” case and make clear the “RESEND Length=-1” mechanism as this seems like the only case where that should be relevant. Likely also a state machine since the state flows are different than the regular data loss mechanism; Tx DATA, Tx Lost, Tx RESEND, Rx RPC_UNKNOWN, Tx DATA or something of the like.
Also as I mentioned in another reply [1] the described mechanism has many undesirable propertys.
2. My statement was more about explicit consideration as to how encryption integrates. Sure, you can just layer it below with IPSec or above with TLS, or any other myriad of ways as you could in literally any protocol for transmitting bytes, but TLS over TCP versus QUIC shows how much more painful it is if it was not in the original design.
3. I did not see that it was a duplex benchmark, but that is a very non-standard way to benchmark, so I can not comment as to comparability to standard one-way benchmarks. If a one-way benchmark is equivalent to the sum of the two-way flows, then I will concede that it is comparable to QUIC implementations.
However, my personal experience with designing and implementing transport protocols results in me finding 8 Gbit/s per core in software to be extremely slow which also applies to every QUIC implementation I have ever seen.
4. Congestion control in QUIC, like in TCP, is not fixed though there are recommended implementations. Furthermore, stream prioritization, something not present in TCP due to the single stream nature, is also not fixed. A implementation that offers prioritizing which streams to transmit data for by SRPT would likely see similar outcomes to Homa. I am not aware of any implementation that does so, but it is a freely configurable element of that protocol. I see nothing that would materially stop a comparable solution from being dropped in.
Mainly, unlike TCP, the design of QUIC has knobs in many of the right places so that you can reasonably configure the desired propertys with minimal design clashing.
Really, the main fundamental distinction in steady state operation is the resend mechanism being largely ack driven TCP and QUIC whereas Homa is largely nack driven.
I am of the opinion that ack driven is much superior technically, which is part of why I think Homa is a bad design, but likely for different reasons and considerations as to how it relates to transport protocol design than the designers of QUIC.
This is also not to say that I think QUIC is the best protocol, QUIC has a number of design and performance weaknesses, but it is overall a reasonable and material improvement over TCP and should be considered the standard to beat.
[1] https://news.ycombinator.com/item?id=49960483
However, this would create a conflict within QUIC. My understanding is that one of QUIC's primary use cases is WAN connections between browsers and Web servers. In this world Homa's congestion control mechanism does not work; it only works in datacenter environments with speed-of-light delays at most a few 10's of microseconds. Thus, QUIC would need to support multiple congestion control mechanisms, one for WAN and one for datacenter and choose between them as appropriate.
In addition, my understanding is that QUIC is stream-based, which would mean, as with TCP, the protocol doesn't know about message lengths. This makes it hard to prioritize short messages.
And, by the way, as much as you may prefer acks to grants, the grant approach is pretty fundamental to Homa's congestion control.
To me, QUIC and Homa were designed for very different environments. Homa will never compete with QUIC in the WAN, and I think it will be hard for QUIC to compete with Homa in the datacenter (Homa can specialize itself entirely for datacenter apps, whereas QUIC can't).
1. Having multiple congestion control algorithms is already standard in both QUIC and TCP. As deploying Homa requires you to control both ends or at least understand both ends, in a comparable QUIC deployment you would know they are in a datacenter and can just link the corresponding implementation. Or you could use the configuration parameter handshake mechanism to switch at runtime for certain flows.
2. QUIC is stream-based, but they support fixed/final-sized streams which are logically equivalent to messages. You just generate a disposable fixed-size stream for each “message”. So, it is trivial for the Tx side to know the “message” size and thus do SRPT prioritization. The thing that would be more annoying to support is the Rx side knowing the stream size as size finalization is detected by the Rx side via the FIN bit when transmitted the terminal segment in a stream. You could transmit the terminal segment first, but that is awkward and Homa is superior in that respect. However, Rx size finalization would only really be important for avoiding over-allocation of Rx buffering in QUIC since they do not rely on MAX_STREAM_DATA for flow control whereas the equivalent concept in Homa, GRANT, is used for flow control and thus precise Tx sizing awareness by the Rx side is more important.
3. As mentioned above, Homa GRANT is very similar to QUIC MAX_STREAM_DATA which are basically just extensions of TCP window size to more than one data flow. The main difference is how they are used.
TCP and Homa both use it for Rx resource availability and bandwidth limiting. Homa uses it additionally for data flow prioritization as it has more than one data flow. QUIC uses it for Rx resource availability and implementations can use it for data flow prioritization. QUIC uses MAX_DATA or just plain congestion control bandwidth estimation for bandwidth limiting.
The fundamental usage is Rx resource availability. If you have no Rx buffers available to take the data, then no amount of data sent is meaningful. That is a hard cap. However, using that as a bandwidth limiting mechanism like TCP and Homa requires continuous tight feedback to avoid the Tx from determining there are no Rx resources available while not over-provisioning allowed bandwidth.
To illustrate, suppose I have a 10 GB/s NIC and generate a 1 GB request per second, but my end-to-end is a 1 GB/s link with 1 ms one-way latency. So, this link only handles 1 MB in-flight without buffering. For the sake of simplicity say we pack only 1 KB per packet, so the link can only have 1,000 packets in flight without buffering. How should we manage the GRANT flow?
If we do a initial GRANT of 10 MB, then the Tx will shovel 10 MB into its 10 GB/s NIC which will spit out into the network in 1 ms. This will clog up at the 1 GB/s bottleneck requiring ~10 MB of buffering at the bottleneck. The Tx will then wait for more GRANT as it just bursted it all. After 1 ms the Rx will get the first byte, what should it GRANT and when? If we wait 1 ms, we have received our first MB, should we GRANT 1 MB? Our steady state then looks like bursts of 10 GB/s with 10 MB of buffering at the bottleneck until the request completes.
If we do a initial GRANT of 1 MB then GRANT 1 MB when we get it all, then we stall.
If we do a initial GRANT of 2 MB, then 1 MB every ms then we get 2 MB of buffering at the bottleneck until the request completes.
What is the optimal flow? The optimal outcome is a continuous 1 GB/s flow with one 1 KB packet entering every 1 us. This results in minimal buffering as a result of the transport protocol. For the GRANT mechanism to do this would require a new 1 KB GRANT to arrive every 1 us.
Absent a separate point-to-point bandwidth control system to pace the DATA flow you are not getting that. But once you have that, the GRANT is no longer the bandwidth limiter, it becomes a Rx resource availability and flow prioritization mechanism, just like QUIC. It no longer becomes meaningfully distinct.
However, software QUIC implementations are generally much slower than software TCP implementations. This is not due to any fundamental protocol design limitations, it is just that most/every QUIC implementation is poorly implemented for maximizing performance.
You can, of course, do multiple times more throughput than QUIC or TCP with a protocol better designed for performance, but that is likely orthogonal to your question.
bold claim. Go tell that to Cloudflare, for example. And/or show us your implementation.
The official MsQuic benchmarks [2] tuned for the standard QUIC benchmark [3] with massively over-provisioned hardware does not even reach 8 Gbit/s on large sends with negligible RTT and loss.
If the highest performing public implementation can not even reach a paltry 8 Gbit/s it is fair to say they are all poorly performing.
As for what I would consider “okay” that would be around 30 Gbit/s per core, bottlenecking on decryption (assuming you do not use the allowed null encryption for benchmarking), or bottlenecking on ~1/3 of the bulk 1 Kbyte memcpy() throughput on your hardware, whatever comes first. “Good” is on the order of 100 Gbit/s per core or the same external limiting factors.
[1] https://www.sciencedirect.com/science/article/pii/S014036642...
[2] https://microsoft.github.io/msquic/
[3] https://datatracker.ietf.org/doc/html/draft-banks-quic-perfo...
2. TCP doesn't have encryption either. They should probably show DTLS working though.
4. You can't just say QUIC is faster without testing it.
That is a sender-driven resend and as such relies on the sender identifying the “send lost” condition. However, the entire protocol is designed around not doing that and thus you can only feasibly rely on “should have received a reply by now” as your timeout.
But Homa is intended to be a RPC protocol. So you send to server, wait for server to process command, then wait for server to send reply. Your timeout depends on the variable and heterogeneous command processing time.
Even if you were able to give a separate correctly tuned timeout for every possible RPC that is still awful. Any RPC with long processing time should not trigger the timeout until the expected reply time, but if that is far larger than the RTT then you are waiting a tremendous amount of time.
For instance, a select query that is only a few bytes long (and thus fits in one packet) could take seconds on a large database even though the database server is physically nearby and only microseconds away. In that case you would have to wait seconds before timing out instead of just microseconds like a ack-based design could achieve. The worst-case transport latency becomes application-specific instead of related to the application-independent transport parameters like RTT. That is troubling.
2. TCP can be forgiven for that given it predates asymmetric cryptography and even DES. Not considering it relevant in a new protocol this side of the millennium is much less reasonable.
4. MsQuic at 7.5 Gbit/s: https://microsoft.github.io/msquic/
I don't see anything particularly Homa-specific in these issues. A protocol like TCP that doesn't implement RPCs doesn't have to face these issues; it punts them up to the application, so the application has to deal with them. Homa implements all of this in the transport, so applications don't have to worry about it.
In a lossless RPC the client should expect the reply in 1002 ms. When should the client resend the request if it got lost?
Given that the client would not expect to see the first byte of the reply until 1002 ms and receives no other feedback until then, they should not issue a resend of their request until 1002 ms.
That means if the network lost their first request, they will not send a resend until 1002 ms and then get the reply to their resend at 2004 ms.
In contrast, in a ack based approach, the client would expect a server ack in ~2 ms as the server gives feedback on receipt rather than on processing completion/reply start.
That means if the network lost their first request they will not send a resend until 2 ms and then get the reply to their resend at 1004 ms, nearly half the latency.
In ack based approaches, packet loss requires one additional RTT latency to recover per repeated packet loss. But Homa incurs one additional RTT + processing latency to recover per repeated packet loss.
As a secondary downside, this requires the client to hold Tx resources for RTT + processing time since the client does not know when the server has received all of the data. In ack-based approaches, you learn the server has received all of the request when it acks all of the data in the request which is on the order of the RTT and thus can release Tx resources at that time. This is similar to the reason Homa has ACK/NEED_ACK packets except applied to client side resources instead of just server side resources.
The case you are describing is more like what if the server explodes, but you still need the results/modifications related to your RPC to occur. In that case, whether the Tx got through or not is irrelevant because your RPC still did not complete and you need it to complete, so there is no advantage to the ack-based approach. But by that logic, what if it explodes before persisting changes? Or what if all your servers explode? At some point we should call that out of scope for a transport protocol.
I prefer setting that line at “bytes got from one end to the other”. If you want higher guarantees then you can layer that. As you point out, it is easy to just layer a “server dead” timeout ping which basically solves the problem in the same way Homa needs to.
And then in basically every other respect it is worse including compatibility.
Because the NSA and similar organizations will 100% infiltrate the physical infrastructure. It’s easier and harder if not impossible to detect.
This is self-evident, isn't it?
Other than the buzzword factor, why do you think name-dropping AI makes a difference?
On the other hand, if the traffic follows a very predictable pattern, a custom implementation can be much more efficient. Specially nowadays machine learning can find much better solutions through reinforcement learning. And AI cluster data flow is much more predictable than what goes over the internet as a whole.
https://arxiv.org/abs/2304.01433
Basically they use lens array (one lens per input fiber) to focus each specific beam on its assigned mirror and mirror directs the light to the chosen output that also each has a lens in front of the fiber going out.
1 - https://news.ycombinator.com/item?id=48287718 is my post on tla formal modeling complete with examples and decent pdf to go with it. Self contained assumes no background on tla
2 - part of the reason for doing (1) was to model homa's basic protocol. See pdf. I did not go into homa details there because I wanted to stay focused on tla.
3 - I have a kernel bypass homa transport api in development
4 - as wonderful as homa is for intra-data-center IO, the srpt part of the protocol is a bit tough to performantly code. More or less homa wants each NIC resource to be managed by one owner who chooses from among all tx streams that could be sent, which one has smallest srpt (shortest remaining processing time). This applies despite the fact client/servers sharing the NIC run in their own threads. Ive been using mpsc memory rings moving the smallest amount of data i can.
> Link Level Flow Control: InfiniBand uses a credit-based algorithm to guarantee lossless HCA-to-HCA communication. RoCE runs on top of Ethernet. Implementations may require lossless Ethernet network for reaching to performance characteristics similar to InfiniBand. Lossless Ethernet is typically configured via Ethernet flow control or priority flow control (PFC). Configuring a Data center bridging (DCB) Ethernet network can be more complex than configuring an InfiniBand network.[19]
* https://en.wikipedia.org/wiki/RDMA_over_Converged_Ethernet
> A sending station (computer or network switch) may be transmitting data faster than the other end of the link can accept it. Using flow control, the receiving station can signal the sender requesting suspension of transmissions until the receiver catches up. Flow control on Ethernet can be implemented at the data link layer.
* https://en.wikipedia.org/wiki/Ethernet_flow_control
Consider a very simple topology:
Say hosts A and B are both sending data to C, as fast as they can, via switches S1 and S2 (which are connected via a high-speed link). And say the sum of these two flows is more than the capacity of the link to C.S2 is receiving packets destined for C faster than it can forward them, but sending an Ethernet pause frame from S2 to S1 is not a very productive way to alleviate the situation, because it also disrupts any traffic that would be bound for D. It just moves the bottleneck elsewhere and causes collateral damage.
I agree that Ethernet flow control is insufficient, but given that NVidia in an act of brilliant foresight acquired Mellanox a decade ago, I’m fairly certain that this is how all these AI clusters are actually deployed, not using TCP, and maybe not even using Ethernet but infiniband instead.
if not, you get RoCE over ethernet.
Worth noting that Nvidia/Mellanox is currently the only vendor for IB, but that has not always been the case. I had Intel IB switches for the backend network of my Isilon at my last job (we starting using them when they were 'just' Isilon, then EMC Isilon, and now Dell-EMC Isilon); only post-Dell did they start going Ethernet-only, and now with AI/ML they're re-offering IB as an option (LOL).
So Nvidia is kind of being 'rewarded' for sticking with a technology that everyone gave up on (I guess HPC wasn't big enough to bother with).
Besides IB, the other options being used in Top 500 are Slingshot (Eth-base?) and Omni-path (Category: Interconnect Family):
* https://www.top500.org/statistics/list/
For IB the answer is generally a centralized controller, a lossless flow control algorithm and the credit system with subnet coordination. Ie. the packets that would be dropped due to buffer problems never get sent because the originating node knows it's out of credits. This means packet drops should be nonexistent and so the protocols make the assumption only one packet in a zillion ever gets dropped. On the down side: the tiniest of problems on IB networks have incredible impact (but it's realistic to just not have those)
IB doesn't interface with others. Yes, there's an ethernet bridge but it's not going to work nearly as well once you insert that bit of kit.
In ethernet packets are dropped when congestion is experienced, and you configure through priority queues whose packets get dropped first (you configure priority), and "stuff will sort itself out" like on the internet. RoCE tries to construct IB semantics from this, but there's limits. On the plus side: this is how the internet works and so the whole actually works somewhat reasonably over internet or other long-distance links (and tunnels into Neoclouds and/or Clouds). On the downside it's never going to match IB in speed (it's still going to get to 90% or so link speed, 98% if you just naively check the link utilization and ignore the effect retransmissions will have)
When a link becomes saturated working out how to manage that is a hard problem.
I have only used RoCE once at scale, it was really finicky. We would get big waves to pause frames that stalled everything.
That is relevant for musks in-space computation. There, latency between compute nodes, due to being hundreds of kilometres apart, cannot be low.
This pretty much limits the size of AI models which can be trained on a space datacenter to a single satellite to avoid incurring huge latency overhead both in training and inference.
With the smartest models getting bigger and bigger, that really doesn't look good for space based compute.
Blog post: https://openai.com/index/mrc-supercomputer-networking/
Paper: https://cdn.openai.com/pdf/resilient-ai-supercomputer-networ...
Spec: https://www.opencompute.org/documents/ocp-mrc-1-0-pdf
MRC is an extension of RoCE that does out-of-order packet delivery, packet spraying across all available paths using SRv6 routing, combination of ECN and packet trimming for congestion control (so sender-based CC like TCP and unlike Homa).
> Today's show puts the “heavy” in Heavy Networking with a discussion about Multipath Reliable Connection (MRC). MRC is about getting an even distribution of RDMA traffic across equal cost multipath links in massive data centers, and doing so over plain old lossy Ethernet. Ethan, Drew, and guests discuss how package spraying, resilience to link and fabric failures, and congestion control allow MRC to be a specialized transport for AI workloads.
> Our guests today are from Broadcom: Rip Sohan, Distinguished Engineer; Eric Davis, Master Engineer; and Eric Spada, Technical Director and Distinguished Engineer.
* https://www.youtube.com/watch?v=AQhoTd9rf60
https://www.usenix.org/system/files/atc21-ousterhout.pdf