From 0dc6ebd5259133f9b594a67cd5d8884483b0a779 Mon Sep 17 00:00:00 2001 From: Ettore Di Giacinto Date: Tue, 1 Sep 2026 23:30:10 +0000 Subject: [PATCH] fix(cluster): stop blaming a peer for the caller's own expired deadline Review round 1 on the end-to-end proof. Zero blocking items, eleven non-blocking, and three of them turned out to be production defects rather than notes on the report. The one that matters is a misclassification the phase is built to prevent. A dial carries the caller's deadline down to the socket, so when the budget runs out the socket's timer fires and the error travels back up through the WebSocket handshake and the multiplexer. The context's cancellation is a separate timer whose func the scheduler has to run before ctx.Err() stops returning nil, and nothing orders the two. Under contention the socket's error is back in PeerPool.Open first, ctx.Err() reads nil, and a peer that is listening and healthy is reported as ErrPeerUnreachable to a caller that simply ran out of time. An unreachable peer is a fact a caller may act on and an expired deadline is not, and core/services/nodes routes around a replica it is told is unreachable. callerRanOut answers that question in one place: ctx.Err() when it is set, and otherwise the wall clock against the caller's own deadline. That is sound because it is the same instant the socket compared itself against, so if the socket's timer fired this comparison is past it too. The ambiguous instant resolves towards the caller, which is the direction that never blames a peer. The spec that caught it, peerlink_test.go's "blames the caller's deadline", was red in three of seven -race runs and had been since Task 5, which is often enough to read as noise and is why single-run verification never saw it. Rather than leave the proof to a coin flip, a second spec makes the window deterministic: Open is handed a context whose deadline has passed and whose cancellation has not been delivered, against an address nothing is listening on, so the dial fails for real. It reddens without the fix. The peer link's yamux windows were applied to one end only. A receive window is advertised by the side that RECEIVES, so configuring the dialler alone tunes exactly one direction, and the direction left on the 256 KiB default is the one that carries a relayed model artifact INTO the replica that owns the worker's tunnel. That is the largest thing the link ever moves and it is the direction the load measurement exercises: the review read it as flowing toward the dialler and it does not. PeerLinkConfig is now exported and used on both ends. Measured, same box, 128 MiB staged through the relay against the same transfer without one: the relayed path cost 1.6x to 2.0x the direct path's transfer window before, and 1.06x to 1.25x after. The SSRF reachability spec could be fooled into reporting an SSRF that did not happen. It bound the victim on 127.0.0.2 at an ephemeral port and required 127.0.0.1 at the same port to refuse, so any other spec in the run holding that number made the dial succeed; red one run in seven, green five of five in isolation. It now picks from below the kernel's ephemeral range, the same fix the harness got for the adjacent-port collision. The rest are the specs and the report saying what they mean. Scenario 1's advertisement assertion could not tell "the worker advertises nothing" from "the JSON key moved", which matters because removing the advertisement is the change it covers. It was green against a renamed key. The roster now keeps the raw key set beside the decoded fields and the spec requires both keys present before reading them as empty. Scenario 4's refusal-body check was a four-way disjunction admitting bare "tunnel", "not connected" and "unroutable". Those alternatives were inert and each would be satisfied by refusals that say nothing about routing, in the one assertion the whole negative control rests on. It is "no route" alone. The head-of-line gate bounded the worst probe by the whole transfer window, which admits about eightfold degradation and loosens as the box slows. It is now half the window, plus a scale-free ratio against the worst probe under the SAME cold load with nothing to transfer, which is the control that isolates the transfer from the load. Not tighter than that, and the reason is measured rather than cautious: under a concurrent -race suite the worst relayed probe reached a fifth of its window, so a quarter-window gate would have had 1.2x of margin, and a spec that fails one run in three is worse than no spec. The report entry printed p90 and p99 off samples of twenty, where both land on the same element and p99 often lands on the max, so one number appeared three times under three names. A quantile is now printed only when the sample can separate it. Two claims in the report were wrong and are withdrawn rather than softened. Scenario 2's race is closed by the trailing re-read of the owner, not by the pre-assertion the report credited: a move to the non-owner mid-request would serve directly and still return 200, and only the trailing read reddens on it. And "the median request is unchanged" holds on this box and not on the reviewer's, where the relayed median rises up to 82% and p99 up to 3.5x. What survives on both is structural: the worst probe is a small fraction of the window in which bytes are moving, so the session interleaves rather than serialising. Sharing a session with a bulk transfer costs latency; it does not cost service. The disk footprint note undercounted, and the reviewer lost a run to a full disk on this box, so it is worth having right: two bulk models seeded into two frontends and staged to the worker is about 768 MiB, not 512 MiB. Left alone deliberately: the worker's backend port allocator still hands out ports without checking they are free, and its default range still overlaps the kernel's ephemeral range. It is confirmed, it is out of scope here, and it is being tracked as a named follow-up rather than fixed under an e2e task. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto --- core/http/endpoints/cluster/peer.go | 18 +- core/services/cluster/peerlink.go | 63 ++++- core/services/cluster/peerlink_test.go | 41 ++++ core/services/worker/tunnel_test.go | 54 ++++- .../e2e/distributed/cluster_baseline_test.go | 39 +++- tests/e2e/distributed/cluster_tunnel_test.go | 219 +++++++++++++----- 6 files changed, 363 insertions(+), 71 deletions(-) diff --git a/core/http/endpoints/cluster/peer.go b/core/http/endpoints/cluster/peer.go index f0fb96ebc..00ef12c5b 100644 --- a/core/http/endpoints/cluster/peer.go +++ b/core/http/endpoints/cluster/peer.go @@ -52,7 +52,23 @@ func PeerHandler(token string, onSession func(peerID string, sess *yamux.Session // Server side of the mux: the dialing peer is the client, so it owns // the odd stream IDs and this side the even ones. - sess, err := yamux.Server(clustersvc.WebsocketConn(ws), nil, nil) + // + // The SAME configuration the dialler uses, and that is load bearing + // rather than symmetry for its own sake. A yamux receive window is + // advertised by the receiving side, so a nil here left this end on the + // 256 KiB default while the dialler ran at 4 MiB, and the direction + // governed by this end is the one that carries a relayed model artifact + // INTO the replica that owns the worker's tunnel. That direction was + // measured at roughly half the throughput of the same transfer without + // a relay in it. + // + // It also puts the same ceiling on unread data at this end that + // PeerLinkConfig already documents for the dialling end, so a replica + // is now sized against that figure per link in BOTH directions. That + // is the cost of the window being useful at all: a window is a bound + // on data received and not yet read, so a receiver that will not + // buffer cannot advertise one. + sess, err := yamux.Server(clustersvc.WebsocketConn(ws), clustersvc.PeerLinkConfig(), nil) if err != nil { xlog.Error("cluster peer link session setup failed", "peer", peerID, "error", err) _ = ws.Close() diff --git a/core/services/cluster/peerlink.go b/core/services/cluster/peerlink.go index 5f5e31864..fec713891 100644 --- a/core/services/cluster/peerlink.go +++ b/core/services/cluster/peerlink.go @@ -103,8 +103,20 @@ const ( peerLinkMaxWindow = 32 * 1024 * 1024 ) -// peerLinkConfig returns the yamux configuration for a replica-to-replica link. -func peerLinkConfig() *yamux.Config { +// PeerLinkConfig returns the yamux configuration for a replica-to-replica link. +// +// Exported because BOTH ENDS need it and only one of them lives here. A yamux +// receive window is advertised by the side that RECEIVES, so a link configured +// on the dialler alone is tuned in exactly one direction: bytes travelling from +// the accepting replica back to the dialler get these windows, and bytes +// travelling the other way get yamux's defaults. The other way is the one that +// carries a relayed model artifact to the replica that owns the worker's +// tunnel, which is the largest thing this link ever moves. +// +// A fresh config per call, never a shared one: yamux keeps the pointer for the +// life of the session, and two sessions sharing one struct would share whatever +// a future field on it comes to mean. +func PeerLinkConfig() *yamux.Config { cfg := yamux.DefaultConfig() cfg.InitialStreamWindowSize = peerLinkInitialWindow cfg.MaxStreamWindowSize = peerLinkMaxWindow @@ -183,9 +195,9 @@ func (p *PeerPool) Open(ctx context.Context, peerID string) (net.Conn, error) { if err == nil { return st, nil } - // A caller whose own context expired must not cost every other worker + // A caller whose own budget expired must not cost every other worker // its link: the session is fine, this request is not. - if ctxErr := ctx.Err(); ctxErr != nil { + if ctxErr := callerRanOut(ctx); ctxErr != nil { return nil, ctxErr } // A session that died between calls is the common case, not an @@ -202,10 +214,10 @@ func (p *PeerPool) Open(ctx context.Context, peerID string) (net.Conn, error) { // the peer, which may be listening and perfectly healthy. Blaming it // would let one impatient client get a good replica routed around. // - // This also swallows a genuine ErrInstanceNotFound when the context + // This also swallows a genuine ErrInstanceNotFound when the budget // happened to expire at the same moment, which is the safe direction: // a timeout must never be able to manufacture absence. - if ctxErr := ctx.Err(); ctxErr != nil { + if ctxErr := callerRanOut(ctx); ctxErr != nil { return nil, ctxErr } return nil, err @@ -216,7 +228,7 @@ func (p *PeerPool) Open(ctx context.Context, peerID string) (net.Conn, error) { // The peer answered and completed a handshake but will not carry a // stream, which is a transport condition and never absence. _ = sess.Close() - if ctxErr := ctx.Err(); ctxErr != nil { + if ctxErr := callerRanOut(ctx); ctxErr != nil { return nil, ctxErr } return nil, unreachablePeer(peerID, err) @@ -226,6 +238,41 @@ func (p *PeerPool) Open(ctx context.Context, peerID string) (net.Conn, error) { return st, nil } +// callerRanOut reports whether the CALLER's budget is what ended an attempt, +// and is the one place that question is answered. +// +// ctx.Err() alone is not that question, and the difference is a real +// misclassification rather than a nicety. A dial carries the caller's deadline +// down to the socket, so when the budget runs out the socket's own timer fires +// and the error travels back up through the WebSocket handshake and the +// multiplexer. The context's cancellation is a SEPARATE timer whose func has to +// be run by the scheduler before ctx.Err() stops returning nil, and nothing +// orders the two. Under contention the socket's error can be back here first, +// ctx.Err() reads nil, and a peer that is listening and healthy is reported as +// ErrPeerUnreachable to a caller that simply ran out of time. +// +// That is the exact confusion this package refuses everywhere else: an +// unreachable peer is a fact about the peer that a caller may act on, and an +// expired deadline is a fact about the caller that it may not. The wall clock +// settles it without waiting for a goroutine, because the deadline is the same +// instant the socket compared itself against: if the socket's timer fired, this +// comparison is past it too. +// +// A context with no deadline falls through to ctx.Err(), which is the whole +// answer for cancellation: a Canceled context has already had its error set by +// the caller of cancel, with no timer in between. +func callerRanOut(ctx context.Context) error { + if err := ctx.Err(); err != nil { + return err + } + // Not time.Now().After: a failure at exactly the deadline is the caller's + // too, and the ambiguous instant is resolved towards never blaming a peer. + if deadline, ok := ctx.Deadline(); ok && !time.Now().Before(deadline) { + return context.DeadlineExceeded + } + return nil +} + // link returns the per-peer entry, creating it on first use. // // Entries are never pruned: a peer id opened once keeps its entry, and any @@ -300,7 +347,7 @@ func (p *PeerPool) dial(ctx context.Context, peerID string) (*yamux.Session, err // Client side of the mux: the dialling replica owns the odd stream IDs, // matching the yamux.Server the peer handler puts on its end. - sess, err := yamux.Client(WebsocketConn(ws), peerLinkConfig(), nil) + sess, err := yamux.Client(WebsocketConn(ws), PeerLinkConfig(), nil) if err != nil { _ = ws.Close() return nil, unreachablePeer(peerID, err) diff --git a/core/services/cluster/peerlink_test.go b/core/services/cluster/peerlink_test.go index 692f2ee6c..d4406e5e1 100644 --- a/core/services/cluster/peerlink_test.go +++ b/core/services/cluster/peerlink_test.go @@ -30,6 +30,17 @@ func servePeerRoute(e *echo.Echo, token string, onPeer func(string, *yamux.Sessi e.GET(cluster.PeerPath, clusterep.PeerHandler(token, onPeer)) } +// deadlinePassed is a context whose deadline has elapsed and whose +// cancellation has not been delivered, which is the state a caller is in for +// the moment between the two timers that fire at its deadline. Only Deadline is +// overridden: the embedded context supplies a nil Done and a nil Err, which is +// what a context in that window reports. +type deadlinePassed struct{ context.Context } + +func (deadlinePassed) Deadline() (time.Time, bool) { + return time.Now().Add(-time.Millisecond), true +} + var _ = Describe("Peer pool", func() { var ( db *gorm.DB @@ -261,6 +272,36 @@ var _ = Describe("Peer pool", func() { Expect(err).ToNot(MatchError(cluster.ErrInstanceNotFound)) }) + It("blames the caller's deadline even when its cancellation has not landed yet", func() { + // The same rule as the spec above, at the instant that makes it hard. + // + // A dial carries the caller's deadline down to the socket, so the + // socket's timer and the context's cancellation timer fire at the same + // moment and nothing orders them. The socket's error can be back in + // Open before the scheduler has run the context's cancel func, and in + // that window ctx.Err() is nil while the caller's budget is + // unambiguously spent. Reading only ctx.Err() there reports a peer that + // is listening and healthy as unreachable. + // + // That window is real: this spec's sibling above reproduces it under + // `-race` about three runs in seven, which is exactly often enough to + // be dismissed as noise. Here it is made deterministic instead, by + // handing Open a context in precisely that state: deadline passed, + // cancellation not delivered. Nothing is faked about the dial, which + // runs for real against an address nothing is listening on. + refused, err := net.Listen("tcp", "127.0.0.1:0") + Expect(err).ToNot(HaveOccurred()) + addr := refused.Addr().String() + Expect(refused.Close()).To(Succeed()) + Expect(reg.Register(ctx, "peer-refusing", addr, "test")).To(Succeed()) + + _, err = pool.Open(deadlinePassed{ctx}, "peer-refusing") + Expect(err).To(MatchError(context.DeadlineExceeded)) + Expect(err).ToNot(MatchError(cluster.ErrPeerUnreachable), + "the caller's budget was spent before the dial was made; blaming the peer for it is how an impatient client gets a healthy replica routed around") + Expect(err).ToNot(MatchError(cluster.ErrInstanceNotFound)) + }) + It("dials once when many callers open the same peer at the same time", func() { // Without a per-peer lock held across the dial, every concurrent // caller races to dial and all but one of the resulting sessions is diff --git a/core/services/worker/tunnel_test.go b/core/services/worker/tunnel_test.go index b5724a542..8a59f686f 100644 --- a/core/services/worker/tunnel_test.go +++ b/core/services/worker/tunnel_test.go @@ -5,6 +5,7 @@ import ( "encoding/binary" "fmt" "io" + "math/rand/v2" "net" "net/http" "net/http/httptest" @@ -598,6 +599,53 @@ func portOf(ln net.Listener) int { return port } +// listenOnSecondLoopback binds 127.0.0.2 on a port that 127.0.0.1 does not +// have anything on, and will not be given anything on. +// +// The port choice is the assertion's, not an incidental. The spec it serves +// proves a reachability fact: only 127.0.0.2 is listening, so a service that +// honoured the host the frontend named would connect and one that dials +// loopback cannot. A port taken from :0 lands in the kernel's ephemeral range, +// where some unrelated socket on 127.0.0.1 can be holding the same number, and +// then the dial to 127.0.0.1 succeeds and the spec reports an SSRF that did not +// happen. That is not hypothetical: it failed about one run in seven under +// `-race` while passing every time in isolation. +// +// Choosing from BELOW the ephemeral range (32768 on Linux by default) is what +// removes it, because the kernel does not hand those out for outbound +// connections. 127.0.0.1 is probed and released rather than held: holding it +// would make the dial the spec expects to fail succeed instead. +func listenOnSecondLoopback() (net.Listener, int) { + GinkgoHelper() + const ( + floor = 20000 + ceiling = 31000 + attempts = 200 + ) + for i := 0; i < attempts; i++ { + port := floor + rand.IntN(ceiling-floor) + free, err := net.Listen("tcp", fmt.Sprintf("127.0.0.1:%d", port)) + if err != nil { + continue + } + if err := free.Close(); err != nil { + continue + } + victim, err := net.Listen("tcp", fmt.Sprintf("127.0.0.2:%d", port)) + if err != nil { + // A host with no second loopback address fails on every port, so + // this is the skip the spec used to make inline. + if i == 0 && strings.Contains(err.Error(), "assign requested address") { + Skip("this host cannot bind a second loopback address: " + err.Error()) + } + continue + } + return victim, port + } + Fail(fmt.Sprintf("no port in [%d, %d) was free on 127.0.0.1 and bindable on 127.0.0.2 after %d attempts", floor, ceiling, attempts)) + return nil, 0 +} + // The routing table is the security boundary of the whole tunnel, and until now // nothing exercised it: every spec above installs dialLocalTCP, which is exactly // the permissive dialler loopbackService exists to prevent. A review turned @@ -643,12 +691,8 @@ var _ = Describe("Worker tunnel local services", func() { // property of the code. The only listener is on 127.0.0.2; nothing // is on 127.0.0.1 at that port. A service that honoured the named // host would connect; one that dials loopback cannot. - victim, err := net.Listen("tcp", "127.0.0.2:0") - if err != nil { - Skip("this host cannot bind a second loopback address: " + err.Error()) - } + victim, port := listenOnSecondLoopback() DeferCleanup(func() { _ = victim.Close() }) - port := portOf(victim) conn, err := loopbackService(port, port)(ctx, victim.Addr().String()) if err == nil { diff --git a/tests/e2e/distributed/cluster_baseline_test.go b/tests/e2e/distributed/cluster_baseline_test.go index d23ba8f6e..9ef29a976 100644 --- a/tests/e2e/distributed/cluster_baseline_test.go +++ b/tests/e2e/distributed/cluster_baseline_test.go @@ -1,6 +1,7 @@ package distributed_test import ( + "encoding/json" "fmt" "net/http" "os" @@ -43,6 +44,25 @@ type node struct { // have dialled instead of the tunnel. Address string `json:"address"` HTTPAddress string `json:"http_address"` + // keys is what the payload actually carried, which a decoded struct cannot + // tell you. Both fields above are the zero value when a worker advertises + // nothing AND when the key was renamed or dropped, and the whole point of + // the change these specs cover was removing the advertisement, so a rename + // would leave "it advertises nothing" passing for a payload that no longer + // says anything either way. + keys map[string]json.RawMessage `json:"-"` +} + +// UnmarshalJSON decodes the fields above and keeps the raw key set beside them. +func (n *node) UnmarshalJSON(data []byte) error { + // A distinct type, or this method calls itself. + type decoded node + var plain decoded + if err := json.Unmarshal(data, &plain); err != nil { + return err + } + *n = node(plain) + return json.Unmarshal(data, &n.keys) } // requireBinaries reports whether a missing binary must fail the spec instead of @@ -231,16 +251,25 @@ func (p *rosterProbe) idOf(name string) string { } // advertisementOf returns whatever endpoints the roster last reported a node -// advertising, joined for a failure message. Empty means the node published -// none, which is what a worker on this release does. -func (p *rosterProbe) advertisementOf(name string) string { +// advertising, joined for a failure message, and whether the payload carried +// both advertisement keys at all. +// +// The second result is the assertion, not a detail. Removing the advertisement +// is what the change under test did, so "the node advertises nothing" and "the +// keys that would have carried it are gone from the payload" are the two +// outcomes a spec has to keep apart: the first is the feature working, the +// second is the spec having lost its subject and reporting the feature working +// for any node at all, including one that advertises plenty. +func (p *rosterProbe) advertisementOf(name string) (string, bool) { for _, n := range p.lastSeen { if n.Name != name { continue } - return strings.TrimSpace(strings.Join([]string{n.Address, n.HTTPAddress}, " ")) + _, hasAddress := n.keys["address"] + _, hasHTTP := n.keys["http_address"] + return strings.TrimSpace(strings.Join([]string{n.Address, n.HTTPAddress}, " ")), hasAddress && hasHTTP } - return "" + return "", false } // describe is handed to Should as the failure message. Gomega calls a diff --git a/tests/e2e/distributed/cluster_tunnel_test.go b/tests/e2e/distributed/cluster_tunnel_test.go index 13255576c..9301af5df 100644 --- a/tests/e2e/distributed/cluster_tunnel_test.go +++ b/tests/e2e/distributed/cluster_tunnel_test.go @@ -362,7 +362,14 @@ var _ = Describe("Worker tunnel end to end", Label("Distributed"), Label("Cluste // The worker publishes no endpoint of any kind. Without this the // inference below would be satisfied by a frontend that dialled an // advertised address, which is the path this phase removed. - advertised := probe.advertisementOf(c.WorkerName(0)) + // + // Both halves are needed. An empty value alone would also be what a + // spec sees after the keys are renamed or dropped from the payload, + // and that would leave this reporting "advertises nothing" about a + // node it can no longer see the advertisement of at all. + advertised, carriedKeys := probe.advertisementOf(c.WorkerName(0)) + Expect(carriedKeys).To(BeTrue(), + "the roster payload no longer carries the address and http_address keys, so this spec cannot tell a worker that advertises nothing from one it cannot read") Expect(advertised).To(BeEmpty(), "the worker advertised %q, so this spec cannot tell a tunnelled request from a direct dial", advertised) @@ -408,9 +415,19 @@ var _ = Describe("Worker tunnel end to end", Label("Distributed"), Label("Cluste nonOwner := 1 - owner Expect(nonOwner).ToNot(Equal(owner)) - // This is the whole point of the spec, so it is asserted rather than - // arranged: the request is about to go to a replica the database says - // does not hold this worker's tunnel. + // The request is about to go to a replica the database says does not + // hold this worker's tunnel. Read again here rather than inferred from + // the reading above, so a spec that had derived the index some other + // way still could not send it to the owner. + // + // This does NOT close the race, and saying that it does would be the + // same overclaim this phase has had to retract twice: ownership can + // move between this read and the reply, and if it moved TO nonOwner the + // request would be served directly and still come back 200. What closes + // it is the trailing read after the request, which requires the owner to + // be unchanged; a move to nonOwner leaves that read returning nonOwner + // and reddens the spec. This one rules out only the arrangement being + // wrong from the start, which is the cheaper half. Expect(owners.ownerIndexOf(c, 2, nodeID)).ToNot(Equal(nonOwner), "frontend %d owns the worker's tunnel, so a request to it would not be relayed and this spec would prove nothing", nonOwner) @@ -421,10 +438,16 @@ var _ = Describe("Worker tunnel end to end", Label("Distributed"), Label("Cluste expectMockedInference(client, c.FrontendURL(nonOwner), "relayed-model", fmt.Sprintf("frontend %d must relay to frontend %d, which owns the worker's tunnel", nonOwner, owner)) - // The relay did not move ownership. A replica that answered by taking - // the tunnel for itself would also have returned 200 above. + // THIS is the assertion that makes the request above a relayed one. + // + // It rules out the two ways a 200 could arrive without a relay: a + // replica that answered by taking the tunnel for itself, and the + // tunnel moving to nonOwner mid-request so that it served directly. + // Both leave the owner changed, and both redden here. The only window + // left is a move away and back inside one request, which takes two + // claims, and no replica dies in this scenario to prompt either. Expect(owners.ownerIndexOf(c, 2, nodeID)).To(Equal(owner), - "serving a relayed request moved the tunnel to frontend %d, which is not relaying", nonOwner) + "the tunnel is no longer held by frontend %d, so the request to frontend %d was not necessarily relayed", owner, nonOwner) }) // Scenario 3. Kills the replica holding the tunnel. The worker must land on @@ -524,16 +547,20 @@ var _ = Describe("Worker tunnel end to end", Label("Distributed"), Label("Cluste Expect(refused.status).ToNot(Equal(http.StatusOK), "the frontend served an inference for a worker that holds no tunnel, so something other than the tunnel reaches it: %s", refused.body) - // And it fails for the RIGHT reason. A frontend that was refusing for - // any other cause (a missing model, a backend it could not install, an - // unhealthy node) would satisfy the assertion above just as well, and - // would leave the three specs before this one unproven. - Expect(refused.body).To(Or( - ContainSubstring("no route"), - ContainSubstring("not connected"), - ContainSubstring("unroutable"), - ContainSubstring("tunnel"), - ), "the refusal does not name the missing route: %s", refused.body) + // And it fails for the RIGHT reason. A frontend refusing for any other + // cause (a missing model, a backend it could not install, an unhealthy + // node) would satisfy the assertion above just as well, and would leave + // the three specs before this one unproven. + // + // One substring, not a disjunction. This is the strongest leg of the + // whole control, and a disjunction is where such a leg goes soft: the + // looser alternatives this used to carry ("tunnel", "not connected", + // "unroutable") would each be satisfied by refusals that say nothing + // about routing, and one of them is a word this deployment's messages + // are full of. "no route" is what cluster.ErrNoRoute reads as, and + // nothing else on this path produces it. + Expect(refused.body).To(ContainSubstring("no route"), + "the refusal does not name the missing route, so this spec cannot tell a worker with no tunnel from a request that failed for one of the ordinary reasons: %s", refused.body) // The control's own control: put the tunnel back, change nothing else, // and the same request must now succeed. Without this the refusal above @@ -583,13 +610,34 @@ func withBalancer(into **frontendBalancer, arm ...func(*frontendBalancer)) func( } } -// percentile returns the p'th percentile of durations, which must be sorted. -func percentile(sorted []time.Duration, p float64) time.Duration { - if len(sorted) == 0 { - return 0 +// percentileIndex is where the p'th percentile of n sorted samples falls, or +// -1 when there are none. +func percentileIndex(n int, p float64) int { + if n == 0 { + return -1 } - idx := int(float64(len(sorted)-1) * p) - return sorted[idx] + return int(float64(n-1) * p) +} + +// slowestOf is the worst sample, which is the statistic a blocking question +// turns on: a session that stalls one request while a transfer holds it shows +// up in the tail and not in the middle. +func slowestOf(samples []time.Duration) time.Duration { + worst := time.Duration(0) + for _, d := range samples { + if d > worst { + worst = d + } + } + return worst +} + +// sortedCopy returns samples in ascending order without disturbing the caller's +// slice, which the report entry reads again afterwards. +func sortedCopy(samples []time.Duration) []time.Duration { + sorted := append([]time.Duration(nil), samples...) + sort.Slice(sorted, func(i, j int) bool { return sorted[i] < sorted[j] }) + return sorted } // summarise reports the shape of a latency sample. @@ -597,25 +645,50 @@ func summarise(label string, samples []time.Duration) string { if len(samples) == 0 { return label + ": no samples" } - sorted := append([]time.Duration(nil), samples...) - sort.Slice(sorted, func(i, j int) bool { return sorted[i] < sorted[j] }) + sorted := sortedCopy(samples) var total time.Duration for _, d := range sorted { total += d } - return fmt.Sprintf("%s: n=%d mean=%s p50=%s p90=%s p99=%s max=%s", - label, len(sorted), total/time.Duration(len(sorted)), - percentile(sorted, 0.50), percentile(sorted, 0.90), percentile(sorted, 0.99), sorted[len(sorted)-1]) + out := fmt.Sprintf("%s: n=%d mean=%s", label, len(sorted), total/time.Duration(len(sorted))) + // A quantile is printed only when the sample can separate it from the ones + // already printed and from the max. At the sizes here, n around 20 to 50, + // p90 and p99 routinely land on the same element and p99 often lands on the + // last one, and printing one number three times under three names invites a + // reader to compare a tail nothing measured. Omission says "this sample + // cannot answer that"; a repeated number says the opposite. + printed := map[int]bool{percentileIndex(len(sorted), 1): true} + for _, q := range []struct { + name string + p float64 + }{{"p50", 0.50}, {"p90", 0.90}, {"p99", 0.99}} { + idx := percentileIndex(len(sorted), q.p) + if printed[idx] { + continue + } + printed[idx] = true + out += fmt.Sprintf(" %s=%s", q.name, sorted[idx]) + } + return out + fmt.Sprintf(" max=%s", sorted[len(sorted)-1]) } // The deferred question of this phase: what one yamux session does when a large // message and ordinary inference share it, with the relay adding a second hop // for most requests. // -// It is MEASURED here rather than asserted to be fine. Both sides' yamux -// windows are left at their defaults by a deliberate decision, pending these -// numbers, and the numbers are printed as a report entry so a future change to -// those windows has something to be compared against. +// It is MEASURED here rather than asserted to be fine, and the numbers are +// printed as a report entry so a later change has something to be compared +// against. +// +// Be exact about which windows those numbers do and do not speak for. The +// WORKER TUNNEL's two ends both take yamux's defaults (core/services/worker, +// tunnel.go, and core/http/endpoints/cluster/connect.go), and this measurement +// is no reason to change that, but it is also no evidence that they are right: +// a default receive window is limited by the bandwidth-delay product of the +// link, and loopback has no delay to produce one. The PEER LINK's windows are raised +// on both ends. What is measured here is whether a session SERIALISES, which +// loopback answers perfectly well, and not whether a window is large enough for +// a link with latency, which it cannot answer at all. const ( // bulkArtifactSize is the model artifact staged over the tunnel while // probes run. It stands in for the 50MB-class message this feature has to @@ -630,6 +703,12 @@ const ( // hypothetical: the first version of this spec used 64MiB and still passed // with the artifact cut to 4KiB, which is the definition of measuring // nothing. + // + // What it costs, since it is the largest thing this suite puts on disk: + // two bulk models seeded into each of two frontends is 512 MiB, and each is + // then staged to the worker, which is 256 MiB more. About 768 MiB under + // TMPDIR for the length of this spec, plus one 128 MiB string resident in + // the test process. The harness removes the tree in Stop. bulkArtifactSize = 128 << 20 // tinyArtifactSize is the same cold load with nothing to transfer. It is @@ -655,19 +734,55 @@ const ( // reporting. minOverlappingProbes = 5 - // holStallCeiling is the coarse absolute backstop, for the case where the - // bulk transfer is itself pathologically slow and the transfer-window - // comparison would admit a probe latency no deployment would tolerate. + // holStallShare bounds the worst probe as a fraction of the window in which + // bytes were moving. It is the STRUCTURAL assertion: a session that + // head-of-line blocks parks a probe until the transfer lets go, so a + // stalled probe's latency is on the order of the whole window, and one that + // interleaves finishes many probes inside it. // - // The comparison is the assertion that bites; this is only the floor under - // it. A session that head-of-line blocks parks a probe until the transfer - // releases the session, so the stalled probe's latency is on the order of - // the transfer window's; one that interleaves answers a warm mock - // completion in the milliseconds it takes with nothing else on the wire. - // An ABSOLUTE ceiling would have to sit between those two, and on this - // hardware they are 30ms and 250ms, which is too narrow a gap to hold on a - // loaded box. The comparison holds whatever the box's speed, because both - // numbers move with it. + // Half, not the whole window. Bounding by the window itself admits a probe + // that took nearly all of it, which is the wedge with the numbers filed + // off. + // + // Not tighter than half, and the reason is measured rather than cautious. A + // probe's tail grows faster than the transfer window does when the box is + // busy: under a concurrent `-race` suite the worst relayed probe reached + // 20% of its window here, so a quarter would have had 1.2x of margin and a + // spec that fails one run in three is worse than no spec. Half leaves 2.4x + // on the same run and still reddens on a stall, which parks a probe for the + // window rather than a fifth of it. + holStallShare = 2 + + // holStallControlFactor bounds the worst probe against the worst probe + // under the EMPTY load in the same run, which is the second half of not + // relaxing under load: a slower box raises the control and the bound with + // it, while an absolute number would simply admit more. + // + // The empty load is the right thing to compare against and the plain + // baseline is not. Both samples then contain a cold load's contention for + // the worker, the router and the session, and the only thing that differs + // between them is 128 MiB crossing the wire. Compared against the quiet + // baseline instead, a transfer that cost nothing at all would still look + // like a regression on any box where a cold load is expensive. + // + // Eight, from both ends of the gap it has to sit in. Healthy runs measured + // 1.6x to 3.8x on this box under a concurrent `-race` suite, and about 2.5x + // to 3x on the reviewer's; a session that stalled a probe until the + // transfer let go would show the whole window over the same control, which + // is 14x to 43x on the same runs. + holStallControlFactor = 8 + + // holStallCeiling is the coarse absolute backstop under both of those, for + // a transfer so slow that a quarter of its window is a latency no + // deployment would tolerate. + // + // It is deliberately far above anything measured rather than tuned, because + // an absolute number cannot separate a wedge from a slow box. Measured + // worst probe and transfer window, for scale: 21-36ms against 261-576ms on + // the box this was written on, and 70-136ms against 590-1320ms on the + // reviewer's. An absolute ceiling that bit on the first machine's wedge + // would fail on the second machine's healthy run, which is why the two + // relative bounds above are the assertions and this is only a floor. holStallCeiling = 5 * time.Second ) @@ -825,15 +940,15 @@ var _ = Describe("Worker tunnel under load", Label("Distributed"), Label("Cluste report = append(report, line) GinkgoWriter.Println(line) - slowest := time.Duration(0) - for _, d := range underBulk { - if d > slowest { - slowest = d - } - } - Expect(slowest).To(BeNumerically("<", transferWindow), - "%s: a probe waited %s while %s of bytes were moving, which is the shape of a session that stalled the probe until the transfer let go, not of one that interleaved them", + slowest := slowestOf(underBulk) + Expect(slowest).To(BeNumerically("<", transferWindow/holStallShare), + "%s: a probe waited %s of the %s in which bytes were moving, which is the shape of a session that stalled the probe until the transfer let go, not of one that interleaved them", label, slowest, transferWindow) + control := slowestOf(underTiny) + Expect(control).To(BeNumerically(">", 0), "%s: the empty-load control produced no samples", label) + Expect(slowest).To(BeNumerically("<", holStallControlFactor*control), + "%s: the worst probe was %s while bytes were moving against %s under the same cold load with nothing to move, which is a stall rather than the contention a shared session costs", + label, slowest, control) Expect(slowest).To(BeNumerically("<", holStallCeiling), "%s: a probe waited %s while the bulk transfer held the session", label, slowest) }