From cf6b1a9ed74fc2633de46b323b93b96ca417fbca Mon Sep 17 00:00:00 2001 From: "allen.wu" Date: Mon, 21 Sep 2026 16:05:52 +0800 Subject: [PATCH 1/2] feat(node): add multi-endpoint failover for the L1 RPC client --l1.rpc now accepts a comma-separated endpoint list, primary first. Failover lives in an http.RoundTripper below the rpc layer, so every existing consumer keeps taking a plain *ethclient.Client and the rollup binding still receives a full bind.ContractBackend, including the SubscribeFilterLogs method that an HTTP-only client could not honestly implement. An endpoint is judged failed on a transport error, a 408/429/5xx status, or a reply that is not JSON. The last rule is not theoretical: on a network whose resolver answers for names that do not exist, a vanished endpoint returns HTTP 200 with an HTML block page. Judged on status alone that reads as success, and a devnet run showed the node keep using the dead endpoint until every L1 read failed upstream with "invalid character '<'", silently stopping block production with no error and no restart. Each attempt also carries GetBody so net/http can re-send a request whose pooled connection died. Without it that routine keep-alive race surfaces as an endpoint failure and the sticky switch demotes a healthy primary for the life of the process. A single endpoint keeps the previous ethclient.Dial behaviour, including ws:// and IPC, so existing configurations are unaffected. AI-Developer: claude-opus-4.8 AI-Agent: claude code AI-Reviewer: none Harness-Skill: none Co-Authored-By: Claude Opus 4.8 --- node/cmd/node/main.go | 17 +- node/flags/flags.go | 2 +- node/rpcfailover/failover.go | 388 ++++++++++++++++++ node/rpcfailover/failover_test.go | 628 ++++++++++++++++++++++++++++++ 4 files changed, 1029 insertions(+), 6 deletions(-) create mode 100644 node/rpcfailover/failover.go create mode 100644 node/rpcfailover/failover_test.go diff --git a/node/cmd/node/main.go b/node/cmd/node/main.go index 235fedc2d..3513aacf5 100644 --- a/node/cmd/node/main.go +++ b/node/cmd/node/main.go @@ -33,6 +33,7 @@ import ( "morph-l2/node/flags" "morph-l2/node/hakeeper" "morph-l2/node/l1sequencer" + "morph-l2/node/rpcfailover" "morph-l2/node/sequencer" "morph-l2/node/sequencer/mock" "morph-l2/node/sync" @@ -106,15 +107,21 @@ func L2NodeMain(ctx *cli.Context) error { } // ========== Shared L1 client ========== - // One ethclient.Dial per process — all L1-touching components (syncer, - // derivation, l1sequencer Tracker/Verifier/Signer, rollup binding) share - // the same connection pool, retry policy, and metrics surface. Adding a - // new consumer means injecting this client, not opening a new one. + // One dial per process — all L1-touching components (syncer, derivation, + // l1sequencer Tracker/Verifier/Signer, rollup binding) share the same + // connection pool, retry policy, and metrics surface. Adding a new consumer + // means injecting this client, not opening a new one. + // + // --l1.rpc accepts a comma-separated list of endpoints, primary first. When + // more than one is given, failover happens inside the HTTP transport, so + // every consumer below keeps receiving a plain *ethclient.Client. See + // node/rpcfailover for what does and does not count as endpoint failure — + // notably, it assumes the caller only reads, which holds here. l1RPC := ctx.GlobalString(flags.L1NodeAddr.Name) if l1RPC == "" { return fmt.Errorf("%s is required", flags.L1NodeAddr.Name) } - l1Client, err := ethclient.Dial(l1RPC) + l1Client, err := rpcfailover.Dial(context.Background(), "L1", l1RPC, nodeConfig.Logger) if err != nil { return fmt.Errorf("dial l1 node error: %v", err) } diff --git a/node/flags/flags.go b/node/flags/flags.go index e6ed45f78..8168251f2 100644 --- a/node/flags/flags.go +++ b/node/flags/flags.go @@ -60,7 +60,7 @@ var ( L1NodeAddr = cli.StringFlag{ Name: "l1.rpc", - Usage: "Address of L1 User JSON-RPC endpoint to use (eth namespace required)", + Usage: "Address of L1 User JSON-RPC endpoint to use (eth namespace required). Accepts a comma-separated list, primary first, to fail over when an endpoint stops answering; multiple endpoints must all be http(s)", EnvVar: prefixEnvVar("L1_ETH_RPC"), } diff --git a/node/rpcfailover/failover.go b/node/rpcfailover/failover.go new file mode 100644 index 000000000..c74265afb --- /dev/null +++ b/node/rpcfailover/failover.go @@ -0,0 +1,388 @@ +// Package rpcfailover dials an Ethereum JSON-RPC endpoint with transparent +// failover across several endpoints. +// +// Nothing here is specific to L1, or to any particular consumer: the transport +// swaps a target URL and judges a response, and the caller supplies a name used +// only to label logs and errors. To get failover for a different set of +// endpoints, pass a different name — do not reimplement this. +// +// # Constraints +// +// These are real limits of the approach, not accidents of the current caller: +// +// - HTTP(S) only. Over ws:// or IPC the rpc.Client holds a live connection and +// reconnects through its own reconnectFunc, so moving to another endpoint is +// not a matter of rewriting a URL. Dial therefore only engages failover when +// more than one endpoint is configured, and requires http(s) in that case. +// - JSON-RPC only. Health is partly judged by the body starting like a JSON +// document. A REST API — the beacon node, say — needs its own scheme; see +// derivation.FallbackBeaconClient. +// - Read-only, or writes that are idempotent. A failed attempt is re-sent to +// the next endpoint, and "no usable response" does not prove "not executed": +// a proxy can return 502 after the request was processed. Re-sending a +// pre-signed raw transaction is harmless (same hash), but anything order- or +// nonce-dependent needs thought before being wired to this. Note that +// tx-submitter's L2Clients deliberately routes writes to one endpoint only. +// +// # Why the transport layer +// +// Failover lives in an http.RoundTripper rather than in a wrapper around +// ethclient, and that placement is the whole point. rpc/http.go builds a fresh +// http.Request per call, so swapping the target URL underneath the rpc layer is +// invisible to everything above it. Consumers keep taking *ethclient.Client, and +// a contract binding keeps receiving a full bind.ContractBackend — including the +// SubscribeFilterLogs method that no HTTP-only client could honestly implement. +// +// # What counts as a failed endpoint +// +// A transport failure (connection refused, TLS and DNS errors, per-attempt +// timeout), a status of 408, 429 or 5xx, or a reply that is not JSON at all. +// +// That last rule is not theoretical. On a network whose resolver answers for +// names that do not exist, a vanished endpoint does not produce "connection +// refused": the name resolves to a captive-portal address that serves HTTP 200 +// and an HTML block page. Judged on status alone that reads as success, so the +// caller keeps using a dead endpoint and every read fails one layer up with +// "invalid character '<' looking for beginning of value" — no failover, and no +// error the caller can act on. The same shape occurs behind load-balancer error +// pages and misrouted DNS. +// +// The check is the first non-whitespace byte of the body, not the Content-Type +// header: a JSON-RPC reply always starts with '{' or '[', whereas the header is +// sometimes mislabelled (Go's own content sniffer reports text/plain for a JSON +// body), and rejecting a working endpoint is worse than missing a broken one. +// +// Two cases are deliberately out of scope: +// +// - A well-formed JSON-RPC error inside an HTTP 200. That is the endpoint +// doing its job; the error belongs to the caller. Treating it as a failure +// would multiply load across every endpoint and return the same error. +// - An endpoint that answers promptly with a stale view of the chain. Nothing +// at the HTTP layer can see that. +// +// For the L1 caller both are backstopped by the L1Tracker halt state machine, +// which stops block production once the node's L1 view has been blind for too +// long. A new caller needs its own answer for them. +package rpcfailover + +import ( + "bytes" + "context" + "errors" + "fmt" + "io" + "net" + "net/http" + "net/url" + "strings" + "sync/atomic" + "time" + + "github.com/morph-l2/go-ethereum/ethclient" + "github.com/morph-l2/go-ethereum/rpc" + tmlog "github.com/tendermint/tendermint/libs/log" +) + +const ( + // dialTimeout bounds TCP+TLS setup per attempt, so an endpoint whose host is + // blackholed cannot consume the caller's entire deadline. + dialTimeout = 5 * time.Second + + // responseHeaderTimeout bounds the wait for response headers per attempt. + // This is what makes failover work against an endpoint that accepts the + // connection and then never answers — the common shape of a hung RPC + // provider. Because the base transport enforces it per RoundTrip, each + // attempt gets its own budget without wrapping the caller's context. + // + // It does not cover a server that sends headers and then stalls the body; + // that case remains bounded only by the caller's context. + responseHeaderTimeout = 10 * time.Second +) + +// Dial builds a client from a comma-separated list of endpoints, in priority +// order with the primary first. +// +// name identifies which set of endpoints these are ("L1", say). It appears in +// log fields and error messages so that a process dialing more than one set can +// be read, and has no effect on behaviour. +// +// A single endpoint is dialed exactly as a plain ethclient.DialContext would, +// including non-HTTP schemes such as ws:// and IPC paths — failover is opt-in by +// configuring more than one endpoint. Two or more must all be http(s), since the +// failover transport is HTTP-only. +func Dial(ctx context.Context, name, raw string, log tmlog.Logger) (*ethclient.Client, error) { + parts := splitEndpoints(raw) + switch len(parts) { + case 0: + return nil, fmt.Errorf("%s: no rpc endpoint configured", name) + case 1: + return ethclient.DialContext(ctx, parts[0]) + } + + endpoints := make([]*url.URL, 0, len(parts)) + redacted := make([]string, 0, len(parts)) + manageAuth := false + for _, part := range parts { + u, err := url.Parse(part) + if err != nil { + return nil, fmt.Errorf("%s: invalid rpc endpoint %s: %w", name, redactEndpoint(part), err) + } + if u.Scheme != "http" && u.Scheme != "https" { + return nil, fmt.Errorf("%s: rpc endpoint %s: failover across multiple endpoints requires http(s), got scheme %q", + name, redactEndpoint(part), u.Scheme) + } + if u.User != nil { + manageAuth = true + } + endpoints = append(endpoints, u) + redacted = append(redacted, redactEndpoint(part)) + } + + transport := &failoverTransport{ + name: name, + base: newBaseTransport(), + endpoints: endpoints, + redacted: redacted, + manageAuth: manageAuth, + log: log, + } + rpcClient, err := rpc.DialOptions(ctx, parts[0], rpc.WithHTTPClient(&http.Client{Transport: transport})) + if err != nil { + return nil, err + } + log.Info("dialed rpc endpoints with failover", "target", name, "endpoints", strings.Join(redacted, ",")) + return ethclient.NewClient(rpcClient), nil +} + +// splitEndpoints splits a comma-separated endpoint list and trims each entry, +// preserving order. A plain single URL yields a one-element slice, so existing +// single-endpoint configs keep working unchanged. Mirrors +// derivation.Config.BeaconRpcList. +func splitEndpoints(raw string) []string { + var parts []string + for _, part := range strings.Split(raw, ",") { + if part = strings.TrimSpace(part); part != "" { + parts = append(parts, part) + } + } + return parts +} + +// redactEndpoint reduces an endpoint to scheme://host for logs and for the +// errors returned to callers. Hosted providers embed credentials everywhere +// except the host — basic-auth userinfo (https://user:pass@host), an API key in +// the path (Infura /v3/ID, Alchemy /v2/KEY), or one in the query (?apikey=) — +// so echoing the raw string would ship those secrets to whatever aggregates the +// node's logs. Host and port are kept: they are not secret, and dropping the +// port would render two endpoints on the same host indistinguishable, which is +// exactly the devnet and local-node case. Mirrors +// derivation.redactBeaconEndpoint. +func redactEndpoint(raw string) string { + parsed, err := url.Parse(raw) + if err != nil || parsed.Scheme == "" || parsed.Host == "" { + return "" + } + return parsed.Scheme + "://" + parsed.Host +} + +func newBaseTransport() *http.Transport { + tr := http.DefaultTransport.(*http.Transport).Clone() + tr.DialContext = (&net.Dialer{Timeout: dialTimeout, KeepAlive: 30 * time.Second}).DialContext + tr.ResponseHeaderTimeout = responseHeaderTimeout + return tr +} + +// failoverTransport sends each JSON-RPC request to the endpoint it is currently +// stuck to, and walks the remaining endpoints in ring order when that one fails. +type failoverTransport struct { + // name labels logs and errors with which set of endpoints these are. + name string + + base http.RoundTripper + endpoints []*url.URL + redacted []string // endpoints[i] reduced to scheme://host, for logs + + // manageAuth is true when at least one endpoint carries userinfo, which is + // the only case where this transport takes ownership of the Authorization + // header. When no endpoint does, a header supplied via rpc.WithHeader is + // left untouched. + manageAuth bool + + // cur is the index of the endpoint that last answered. Sticking to it + // matters: without it every call would pay the dead primary's timeout again. + // Nothing ever moves it back on its own — a recovered primary is picked up + // when the current endpoint fails, or on the next process restart. + cur atomic.Int32 + + log tmlog.Logger +} + +func (t *failoverTransport) RoundTrip(req *http.Request) (*http.Response, error) { + // rpc/http.go sets req.Body but leaves req.GetBody nil, so the body must be + // buffered here to be replayable against the next endpoint. Request bodies + // are small (a method name and a few arguments). Response bodies, which can + // be megabytes, are never buffered — they are handed back untouched. + var body []byte + if req.Body != nil { + var err error + body, err = io.ReadAll(req.Body) + req.Body.Close() + if err != nil { + return nil, fmt.Errorf("%s: read rpc request body: %w", t.name, err) + } + } + + start := int(t.cur.Load()) + var lastErr error + for i := range t.endpoints { + idx := (start + i) % len(t.endpoints) + + resp, err := t.attempt(req, body, idx) + if err == nil { + if idx != start { + t.cur.Store(int32(idx)) + t.log.Error("switched rpc endpoint", "target", t.name, + "from", t.redacted[start], "to", t.redacted[idx], "cause", lastErr) + } + return resp, nil + } + lastErr = err + t.log.Debug("rpc endpoint attempt failed", "target", t.name, "endpoint", t.redacted[idx], "err", err) + + // No budget left for another endpoint: the caller either cancelled or ran + // out of time. Report that rather than an "all endpoints failed" error + // assembled from attempts that never left the process. + if ctxErr := req.Context().Err(); ctxErr != nil { + return nil, err + } + } + return nil, fmt.Errorf("%s: all %d rpc endpoints failed, last error: %w", t.name, len(t.endpoints), lastErr) +} + +// attempt issues the request against endpoints[idx]. A response that is usable +// by the rpc layer is returned as-is; anything else becomes an error so the +// caller can move on. +func (t *failoverTransport) attempt(req *http.Request, body []byte, idx int) (*http.Response, error) { + target := *t.endpoints[idx] // copy: a RoundTripper must not mutate shared state + attempt := req.Clone(req.Context()) + attempt.URL = &target + attempt.Host = "" // let the Host header follow the new URL + attempt.Body = io.NopCloser(bytes.NewReader(body)) + attempt.ContentLength = int64(len(body)) + // GetBody lets net/http re-send this request itself when the connection it + // took from the idle pool turns out to be dead. Its shouldRetryRequest only + // does that for a body it can rewind — the nothingWrittenError branch returns + // `outgoingLength() == 0 || GetBody != nil`, and a POST with a body fails both + // unless this is set. Leaving it nil surfaces a routine keep-alive race as an + // endpoint failure, and the sticky switch then demotes a healthy primary for + // the rest of the process's life. + attempt.GetBody = func() (io.ReadCloser, error) { + return io.NopCloser(bytes.NewReader(body)), nil + } + if t.manageAuth { + applyBasicAuth(attempt, &target) + } + + resp, err := t.base.RoundTrip(attempt) + if err != nil { + return nil, err + } + if shouldFailover(resp.StatusCode) { + resp.Body.Close() + return nil, fmt.Errorf("http %s", resp.Status) + } + if err := requireJSONBody(resp); err != nil { + resp.Body.Close() + return nil, err + } + return resp, nil +} + +// peekLimit is how far requireJSONBody looks for the first non-whitespace byte. +// A JSON-RPC reply has none, so anything beyond a few bytes means the response +// is not one and the window only needs to be generous, not large. +const peekLimit = 8 + +// requireJSONBody reports an endpoint as failed when its body is not JSON. It +// peeks at the first few bytes and hands them back, so the rpc layer still +// decodes the whole response — full bodies are never buffered, since a single +// eth_getLogs reply can be megabytes. +func requireJSONBody(resp *http.Response) error { + // Only a success is expected to carry a JSON-RPC reply. The statuses that + // reach here otherwise are the 3xx that http.Client will follow and the 4xx + // that shouldFailover deliberately passes through; neither promises a JSON + // body, and judging them here would undo that decision. + if resp.StatusCode < 200 || resp.StatusCode > 299 { + return nil + } + + // A compressed body cannot be inspected without decompressing it, which is + // not worth paying for on every call. Go's transport strips Content-Encoding + // when it decompresses itself, so this only applies when a caller set + // Accept-Encoding explicitly via rpc.WithHeader. + if resp.Header.Get("Content-Encoding") != "" { + return nil + } + + peeked := make([]byte, peekLimit) + n, readErr := io.ReadFull(resp.Body, peeked) + peeked = peeked[:n] + atEOF := errors.Is(readErr, io.EOF) || errors.Is(readErr, io.ErrUnexpectedEOF) + + // Restore the body before returning, on every path: the caller closes it even + // when we reject the response. + resp.Body = replayBody(peeked, resp.Body) + + if readErr != nil && !atEOF { + return fmt.Errorf("read response body: %w", readErr) + } + + switch trimmed := bytes.TrimLeft(peeked, " \t\r\n"); { + case len(trimmed) > 0: + if trimmed[0] != '{' && trimmed[0] != '[' { + return fmt.Errorf("endpoint did not return JSON (body starts with %q, content-type %q)", + trimmed[0], resp.Header.Get("Content-Type")) + } + case atEOF: + // Nothing but optional whitespace in the whole body. + return errors.New("endpoint returned an empty body") + default: + // The peek window held only whitespace and there is more to come. Too + // unusual to judge, and rejecting a working endpoint is the worse error. + } + return nil +} + +// replayBody yields the peeked bytes followed by the rest of the original body, +// and closes the original on Close. +func replayBody(peeked []byte, rest io.ReadCloser) io.ReadCloser { + return struct { + io.Reader + io.Closer + }{io.MultiReader(bytes.NewReader(peeked), rest), rest} +} + +// shouldFailover reports whether an HTTP status justifies trying the next +// endpoint. 5xx and 429 are the endpoint's problem, and 408 means it gave up on +// us. The remaining 4xx are deliberately excluded: a 400/401/403/404 describes +// the request or the credentials, the next endpoint would almost certainly +// reproduce it, and moving on silently would turn a misconfiguration into a +// mystery. A JSON-RPC error arrives as 200 and is not examined here. +func shouldFailover(status int) bool { + return status == http.StatusRequestTimeout || + status == http.StatusTooManyRequests || + status >= 500 +} + +// applyBasicAuth re-derives the Authorization header for the endpoint actually +// being dialed. http.Client.send fills it in from the *original* URL's userinfo +// before RoundTrip runs, so after swapping the URL the header would otherwise +// carry the primary endpoint's credentials to a different provider. +func applyBasicAuth(req *http.Request, target *url.URL) { + if target.User == nil { + req.Header.Del("Authorization") + return + } + password, _ := target.User.Password() + req.SetBasicAuth(target.User.Username(), password) +} diff --git a/node/rpcfailover/failover_test.go b/node/rpcfailover/failover_test.go new file mode 100644 index 000000000..1354b215c --- /dev/null +++ b/node/rpcfailover/failover_test.go @@ -0,0 +1,628 @@ +package rpcfailover + +import ( + "context" + "encoding/json" + "errors" + "io" + "net" + "net/http" + "net/http/httptest" + "net/url" + "strings" + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + tmlog "github.com/tendermint/tendermint/libs/log" +) + +const probeBody = `{"jsonrpc":"2.0","id":1,"method":"eth_blockNumber","params":[]}` + +// recorder is a test endpoint that counts requests and remembers what it saw. +type recorder struct { + server *httptest.Server + hits atomic.Int32 + bodies chan string + auths chan string +} + +// newEndpoint starts a test endpoint that replies with the given status. A 200 +// reply carries a valid eth_blockNumber result so the same helper can be driven +// through ethclient. +func newEndpoint(t *testing.T, status int) *recorder { + t.Helper() + r := &recorder{ + bodies: make(chan string, 8), + auths: make(chan string, 8), + } + r.server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + r.hits.Add(1) + body, _ := io.ReadAll(req.Body) + r.bodies <- string(body) + r.auths <- req.Header.Get("Authorization") + + if status != http.StatusOK { + w.WriteHeader(status) + return + } + var call struct { + ID json.RawMessage `json:"id"` + } + _ = json.Unmarshal(body, &call) + if len(call.ID) == 0 { + call.ID = json.RawMessage("1") + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"jsonrpc":"2.0","id":` + string(call.ID) + `,"result":"0x10"}`)) + })) + t.Cleanup(r.server.Close) + return r +} + +// url returns the endpoint URL, optionally with basic-auth userinfo spliced in. +func (r *recorder) url(userinfo string) string { + if userinfo == "" { + return r.server.URL + } + return strings.Replace(r.server.URL, "http://", "http://"+userinfo+"@", 1) +} + +func newTestTransport(t *testing.T, rawURLs ...string) *failoverTransport { + t.Helper() + endpoints := make([]*url.URL, 0, len(rawURLs)) + redacted := make([]string, 0, len(rawURLs)) + manageAuth := false + for _, raw := range rawURLs { + u, err := url.Parse(raw) + require.NoError(t, err) + if u.User != nil { + manageAuth = true + } + endpoints = append(endpoints, u) + redacted = append(redacted, redactEndpoint(raw)) + } + return &failoverTransport{ + name: "L1", + base: newBaseTransport(), + endpoints: endpoints, + redacted: redacted, + manageAuth: manageAuth, + log: tmlog.NewNopLogger(), + } +} + +// post issues a request shaped like the one rpc/http.go builds: a POST whose body +// is wrapped in a NopCloser, which leaves GetBody nil. +func post(t *testing.T, tr http.RoundTripper) (*http.Response, error) { + t.Helper() + req, err := http.NewRequest(http.MethodPost, "http://placeholder.invalid", io.NopCloser(strings.NewReader(probeBody))) + require.NoError(t, err) + require.Nil(t, req.GetBody, "test request must mimic rpc/http.go, which leaves GetBody nil") + req.ContentLength = int64(len(probeBody)) + req.Header.Set("Content-Type", "application/json") + return tr.RoundTrip(req) +} + +func TestFailoverSwitchesOnServerError(t *testing.T) { + primary := newEndpoint(t, http.StatusServiceUnavailable) + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + + assert.Equal(t, http.StatusOK, resp.StatusCode) + assert.Equal(t, int32(1), primary.hits.Load()) + assert.Equal(t, int32(1), secondary.hits.Load()) + + // The buffered request body must be replayed intact onto the second attempt. + assert.Equal(t, probeBody, <-primary.bodies) + assert.Equal(t, probeBody, <-secondary.bodies) + + // The response must reach the caller unread, for the rpc layer to decode. + body, err := io.ReadAll(resp.Body) + require.NoError(t, err) + assert.Contains(t, string(body), `"result":"0x10"`) +} + +func TestFailoverSwitchesOnConnectionRefused(t *testing.T) { + dead := newEndpoint(t, http.StatusOK) + deadURL := dead.url("") + dead.server.Close() // nothing is listening any more + + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, deadURL, secondary.url("")) + + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + + assert.Equal(t, http.StatusOK, resp.StatusCode) + assert.Equal(t, int32(1), secondary.hits.Load()) +} + +func TestFailoverSticksToNewEndpoint(t *testing.T) { + primary := newEndpoint(t, http.StatusInternalServerError) + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + + for i := 0; i < 3; i++ { + resp, err := post(t, tr) + require.NoError(t, err) + resp.Body.Close() + } + + // The dead primary is probed once, on the call that triggered the switch. + // Later calls must not pay its cost again. + assert.Equal(t, int32(1), primary.hits.Load()) + assert.Equal(t, int32(3), secondary.hits.Load()) +} + +func TestFailoverKeepsClientErrors(t *testing.T) { + primary := newEndpoint(t, http.StatusBadRequest) + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + + // A 400 describes the request, not the endpoint. It is surfaced verbatim and + // the secondary is never consulted. + assert.Equal(t, http.StatusBadRequest, resp.StatusCode) + assert.Equal(t, int32(1), primary.hits.Load()) + assert.Equal(t, int32(0), secondary.hits.Load()) +} + +// newRawEndpoint starts an endpoint that replies with a fixed status, content +// type and body, for responses that are not valid JSON-RPC. +func newRawEndpoint(t *testing.T, status int, contentType, body string) *recorder { + t.Helper() + r := &recorder{bodies: make(chan string, 8), auths: make(chan string, 8)} + r.server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + r.hits.Add(1) + b, _ := io.ReadAll(req.Body) + r.bodies <- string(b) + r.auths <- req.Header.Get("Authorization") + if contentType != "" { + w.Header().Set("Content-Type", contentType) + } + w.WriteHeader(status) + _, _ = w.Write([]byte(body)) + })) + t.Cleanup(r.server.Close) + return r +} + +// TestFailoverOnHTMLBlockPage reproduces the failure seen on the devnet: when an +// endpoint's host stops existing, the corporate DNS resolver answers with a +// captive-portal address that serves HTTP 200 and an HTML block page. Judged on +// status alone that looks like success, so the node kept using a dead endpoint +// and every L1 read failed one layer up with "invalid character '<'". +func TestFailoverOnHTMLBlockPage(t *testing.T) { + primary := newRawEndpoint(t, http.StatusOK, "text/html;charset=utf-8", + "\n") + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + + assert.Equal(t, int32(1), primary.hits.Load()) + assert.Equal(t, int32(1), secondary.hits.Load(), "an HTML body must be treated as endpoint failure") + + body, err := io.ReadAll(resp.Body) + require.NoError(t, err) + assert.Contains(t, string(body), `"result":"0x10"`) +} + +// TestFailoverOnEmptyBody covers a 200 with nothing in it, which is not a +// JSON-RPC reply either. +func TestFailoverOnEmptyBody(t *testing.T) { + primary := newRawEndpoint(t, http.StatusOK, "application/json", "") + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + assert.Equal(t, int32(1), secondary.hits.Load()) +} + +// TestNoFailoverOnJSONRPCError is the other half of the contract: when the +// endpoint answers with a well-formed JSON-RPC error, it is doing its job. The +// error belongs to the caller and must not send us to another endpoint, which +// would multiply load and return the same error anyway. +func TestNoFailoverOnJSONRPCError(t *testing.T) { + primary := newRawEndpoint(t, http.StatusOK, "application/json", + `{"jsonrpc":"2.0","id":1,"error":{"code":-32000,"message":"execution reverted"}}`) + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + + assert.Equal(t, int32(1), primary.hits.Load()) + assert.Equal(t, int32(0), secondary.hits.Load(), "a JSON-RPC error is an answer, not an endpoint failure") + body, err := io.ReadAll(resp.Body) + require.NoError(t, err) + assert.Contains(t, string(body), "execution reverted") +} + +// TestNoFailoverOnBatchResponse guards the array form: batch replies are valid +// JSON-RPC and must survive the check. +func TestNoFailoverOnBatchResponse(t *testing.T) { + primary := newRawEndpoint(t, http.StatusOK, "application/json", + `[{"jsonrpc":"2.0","id":1,"result":"0x1"}]`) + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + assert.Equal(t, int32(0), secondary.hits.Load()) + body, err := io.ReadAll(resp.Body) + require.NoError(t, err) + assert.Equal(t, `[{"jsonrpc":"2.0","id":1,"result":"0x1"}]`, string(body), + "the peeked byte must be handed back to the caller intact") +} + +// TestNoFailoverOnLeadingWhitespace allows a body that starts with whitespace. +func TestNoFailoverOnLeadingWhitespace(t *testing.T) { + primary := newRawEndpoint(t, http.StatusOK, "application/json", + "\r\n {\"jsonrpc\":\"2.0\",\"id\":1,\"result\":\"0x2\"}") + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + assert.Equal(t, int32(0), secondary.hits.Load()) + body, err := io.ReadAll(resp.Body) + require.NoError(t, err) + assert.Contains(t, string(body), `"result":"0x2"`) +} + +// TestNoPeekOnCompressedBody: a compressed body cannot be inspected cheaply, so +// it is passed through rather than guessed at. Go's transport strips +// Content-Encoding when it does the decompression itself, so this only applies +// when a caller set Accept-Encoding explicitly via rpc.WithHeader. +func TestNoPeekOnCompressedBody(t *testing.T) { + primary := &recorder{bodies: make(chan string, 8), auths: make(chan string, 8)} + primary.server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + primary.hits.Add(1) + w.Header().Set("Content-Type", "application/json") + w.Header().Set("Content-Encoding", "gzip") + // Deliberately not valid JSON and not valid gzip: the point is that the + // body is never looked at when it announces an encoding. + _, _ = w.Write([]byte("\x1f\x8b\x08 not-json-but-compressed")) + })) + t.Cleanup(primary.server.Close) + + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + + // Go's transport only leaves Content-Encoding on the response when it did not + // add Accept-Encoding itself — otherwise it decompresses and strips the + // header. Setting it explicitly reproduces the one condition under which the + // guard is reachable, which is a caller using rpc.WithHeader. + req, err := http.NewRequest(http.MethodPost, "http://placeholder.invalid", io.NopCloser(strings.NewReader(probeBody))) + require.NoError(t, err) + req.Header.Set("Accept-Encoding", "gzip") + req.ContentLength = int64(len(probeBody)) + + resp, err := tr.RoundTrip(req) + require.NoError(t, err) + defer resp.Body.Close() + require.Equal(t, "gzip", resp.Header.Get("Content-Encoding"), + "precondition: the transport must not have stripped Content-Encoding") + assert.Equal(t, int32(1), primary.hits.Load()) + assert.Equal(t, int32(0), secondary.hits.Load(), "a compressed body must be passed through unexamined") +} + +// TestFailoverOnHTMLBlockPageWithoutContentType pins that the decision is made +// on the body, not the header: the same block page with no Content-Type at all +// must still be rejected. +func TestFailoverOnHTMLBlockPageWithoutContentType(t *testing.T) { + primary := newRawEndpoint(t, http.StatusOK, "", "") + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + assert.Equal(t, int32(1), secondary.hits.Load()) +} + +// TestClientErrorBodyNotJudged pins the scoping decision: the JSON check applies +// to 2xx only, so a 4xx with an empty or HTML body is still returned verbatim +// rather than being turned into a failover. +func TestClientErrorBodyNotJudged(t *testing.T) { + primary := newRawEndpoint(t, http.StatusForbidden, "text/html", "forbidden") + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + assert.Equal(t, http.StatusForbidden, resp.StatusCode) + assert.Equal(t, int32(0), secondary.hits.Load()) +} + +// breakableConn fails writes once the test flips its flag, while reads keep +// delegating to the real connection. That combination is what produces net/http's +// nothingWrittenError: the write fails before a single byte reaches the wire, and +// because the read side never sees a close, readLoop does not evict the +// connection from the idle pool first (transport.go removeIdleConn). +type breakableConn struct { + net.Conn + broken *atomic.Bool +} + +func (c *breakableConn) Write(b []byte) (int, error) { + if c.broken.Load() { + return 0, errors.New("simulated dead pooled connection") + } + return c.Conn.Write(b) +} + +// TestStalePooledConnectionRetriesSameEndpoint pins that a connection dying in +// the pool is not mistaken for the endpoint being down. +// +// net/http can re-send such a request on a fresh connection by itself, but only +// if it can rewind the body: shouldRetryRequest's nothingWrittenError branch +// returns `req.outgoingLength() == 0 || req.GetBody != nil`, and a POST with a +// body and no GetBody fails both. Without GetBody the error surfaces here +// instead, and the sticky failover then demotes a perfectly healthy primary for +// the rest of the process's life. +func TestStalePooledConnectionRetriesSameEndpoint(t *testing.T) { + primary := newEndpoint(t, http.StatusOK) + secondary := newEndpoint(t, http.StatusOK) + + var broken atomic.Bool + var dials atomic.Int32 + dialer := &net.Dialer{Timeout: 5 * time.Second} + base := http.DefaultTransport.(*http.Transport).Clone() + base.DialContext = func(ctx context.Context, network, addr string) (net.Conn, error) { + conn, err := dialer.DialContext(ctx, network, addr) + if err != nil { + return nil, err + } + // Only the first connection — the one that will be reused — is breakable. + // A redial must succeed, which is what the retry depends on. + if dials.Add(1) == 1 { + return &breakableConn{Conn: conn, broken: &broken}, nil + } + return conn, nil + } + + tr := newTestTransport(t, primary.url(""), secondary.url("")) + tr.base = base + + // First call establishes the connection. Draining the body is what returns it + // to the idle pool. + resp, err := post(t, tr) + require.NoError(t, err) + _, err = io.Copy(io.Discard, resp.Body) + require.NoError(t, err) + require.NoError(t, resp.Body.Close()) + require.Equal(t, int32(1), primary.hits.Load()) + + // The pooled connection dies without the client noticing. + broken.Store(true) + + resp, err = post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + + assert.Equal(t, int32(0), secondary.hits.Load(), + "a connection that died in the pool is not an endpoint failure and must not trigger failover") + assert.Equal(t, int32(2), primary.hits.Load(), + "the request should have been re-sent to the same endpoint on a fresh connection") + assert.Equal(t, int32(0), tr.cur.Load(), "the sticky endpoint must still be the primary") + + body, err := io.ReadAll(resp.Body) + require.NoError(t, err) + assert.Contains(t, string(body), `"result":"0x10"`) +} + +// TestAttemptSuppliesGetBody is the narrow contract the retry above depends on. +func TestAttemptSuppliesGetBody(t *testing.T) { + only := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, only.url(""), only.url("")) + + var captured *http.Request + tr.base = roundTripFunc(func(req *http.Request) (*http.Response, error) { + captured = req + return &http.Response{ + StatusCode: http.StatusOK, + Header: http.Header{"Content-Type": {"application/json"}}, + Body: io.NopCloser(strings.NewReader(`{"jsonrpc":"2.0","id":1,"result":"0x1"}`)), + }, nil + }) + + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + + require.NotNil(t, captured) + require.NotNil(t, captured.GetBody, "net/http needs GetBody to rewind a POST body") + rewound, err := captured.GetBody() + require.NoError(t, err) + replayed, err := io.ReadAll(rewound) + require.NoError(t, err) + assert.Equal(t, probeBody, string(replayed), "GetBody must yield the full original body") +} + +type roundTripFunc func(*http.Request) (*http.Response, error) + +func (f roundTripFunc) RoundTrip(req *http.Request) (*http.Response, error) { return f(req) } + +func TestFailoverOnTooManyRequests(t *testing.T) { + primary := newEndpoint(t, http.StatusTooManyRequests) + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + + assert.Equal(t, http.StatusOK, resp.StatusCode) + assert.Equal(t, int32(1), secondary.hits.Load()) +} + +func TestFailoverAllEndpointsFail(t *testing.T) { + primary := newEndpoint(t, http.StatusBadGateway) + secondary := newEndpoint(t, http.StatusServiceUnavailable) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + + resp, err := post(t, tr) + require.Error(t, err) + assert.Nil(t, resp) + assert.Contains(t, err.Error(), "all 2 rpc endpoints failed") + assert.Equal(t, int32(1), primary.hits.Load()) + assert.Equal(t, int32(1), secondary.hits.Load()) +} + +func TestFailoverAppliesPerEndpointBasicAuth(t *testing.T) { + primary := newEndpoint(t, http.StatusServiceUnavailable) + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url("alice:primary-secret"), secondary.url("bob:secondary-secret")) + + // http.Client.send would have stamped the primary's credentials on the + // request before RoundTrip; reproduce that so the test covers the leak. + req, err := http.NewRequest(http.MethodPost, primary.url("alice:primary-secret"), io.NopCloser(strings.NewReader(probeBody))) + require.NoError(t, err) + req.SetBasicAuth("alice", "primary-secret") + resp, err := tr.RoundTrip(req) + require.NoError(t, err) + defer resp.Body.Close() + + primaryAuth, secondaryAuth := <-primary.auths, <-secondary.auths + assert.NotEmpty(t, primaryAuth) + assert.NotEqual(t, primaryAuth, secondaryAuth, "primary credentials must not be forwarded to the secondary endpoint") + + user, pass, ok := (&http.Request{Header: http.Header{"Authorization": {secondaryAuth}}}).BasicAuth() + require.True(t, ok) + assert.Equal(t, "bob", user) + assert.Equal(t, "secondary-secret", pass) +} + +func TestFailoverDropsInheritedAuthWhenEndpointHasNone(t *testing.T) { + primary := newEndpoint(t, http.StatusServiceUnavailable) + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url("alice:primary-secret"), secondary.url("")) + + req, err := http.NewRequest(http.MethodPost, primary.url("alice:primary-secret"), io.NopCloser(strings.NewReader(probeBody))) + require.NoError(t, err) + req.SetBasicAuth("alice", "primary-secret") + resp, err := tr.RoundTrip(req) + require.NoError(t, err) + defer resp.Body.Close() + + <-primary.auths + assert.Empty(t, <-secondary.auths, "an endpoint without userinfo must not receive another endpoint's credentials") +} + +// TestDialFailoverThroughEthclient checks the whole stack: a dead primary, and +// ethclient decoding a real JSON-RPC response served by the secondary. +func TestDialFailoverThroughEthclient(t *testing.T) { + dead := newEndpoint(t, http.StatusOK) + deadURL := dead.url("") + dead.server.Close() + + secondary := newEndpoint(t, http.StatusOK) + + client, err := Dial(context.Background(), "L1", deadURL+" , "+secondary.url(""), tmlog.NewNopLogger()) + require.NoError(t, err) + defer client.Close() + + height, err := client.BlockNumber(context.Background()) + require.NoError(t, err) + assert.Equal(t, uint64(0x10), height) + assert.Equal(t, int32(1), secondary.hits.Load()) +} + +// TestDialSingleEndpoint pins the compatibility promise: one endpoint behaves +// exactly like the previous ethclient.Dial call. +func TestDialSingleEndpoint(t *testing.T) { + only := newEndpoint(t, http.StatusOK) + + client, err := Dial(context.Background(), "L1", only.url(""), tmlog.NewNopLogger()) + require.NoError(t, err) + defer client.Close() + + height, err := client.BlockNumber(context.Background()) + require.NoError(t, err) + assert.Equal(t, uint64(0x10), height) +} + +func TestDialRejectsNonHTTPWithMultipleEndpoints(t *testing.T) { + _, err := Dial(context.Background(), "L1", "http://a.invalid,ws://b.invalid", tmlog.NewNopLogger()) + require.Error(t, err) + assert.Contains(t, err.Error(), "requires http(s)") + // The rejected endpoint must not be echoed with its credentials. + assert.Contains(t, err.Error(), "ws://b.invalid") +} + +func TestDialAllowsNonHTTPSingleEndpoint(t *testing.T) { + // ws:// with a single endpoint keeps the old code path. The dial fails because + // nothing is listening, not because the scheme was rejected. + _, err := Dial(context.Background(), "L1", "ws://127.0.0.1:1/", tmlog.NewNopLogger()) + require.Error(t, err) + assert.NotContains(t, err.Error(), "requires http(s)") +} + +func TestDialNoEndpoint(t *testing.T) { + _, err := Dial(context.Background(), "L1", " , ", tmlog.NewNopLogger()) + require.Error(t, err) + assert.Contains(t, err.Error(), "no rpc endpoint configured") +} + +func TestSplitEndpoints(t *testing.T) { + assert.Equal(t, []string{"http://a"}, splitEndpoints("http://a")) + assert.Equal(t, []string{"http://a", "http://b"}, splitEndpoints(" http://a , http://b ")) + assert.Equal(t, []string{"http://a", "http://b"}, splitEndpoints("http://a,,http://b")) + assert.Nil(t, splitEndpoints("")) +} + +func TestRedactEndpoint(t *testing.T) { + // Credentials hide in userinfo, path and query; all three must be dropped. + assert.Equal(t, "https://rpc.example.com", + redactEndpoint("https://alice:secret@rpc.example.com/v3/abcdef?apikey=xyz")) + assert.Equal(t, "", redactEndpoint("not a url")) + + // The port must survive, otherwise two endpoints on one host — the devnet and + // local-node shape — become indistinguishable in the failover log line. + assert.Equal(t, "http://127.0.0.1:8545", redactEndpoint("http://127.0.0.1:8545")) + assert.NotEqual(t, + redactEndpoint("http://127.0.0.1:8545"), + redactEndpoint("http://127.0.0.1:8546")) +} + +// TestRedactedEndpointsAreDistinct guards the log line itself: two test servers +// differ only by port, so a redaction that dropped it would make every switch +// message read "from=http://127.0.0.1 to=http://127.0.0.1". +func TestRedactedEndpointsAreDistinct(t *testing.T) { + primary := newEndpoint(t, http.StatusServiceUnavailable) + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + + require.Len(t, tr.redacted, 2) + assert.NotEqual(t, tr.redacted[0], tr.redacted[1]) +} + +func TestShouldFailover(t *testing.T) { + for _, status := range []int{408, 429, 500, 502, 503, 504} { + assert.True(t, shouldFailover(status), "status %d should trigger failover", status) + } + for _, status := range []int{200, 201, 301, 400, 401, 403, 404} { + assert.False(t, shouldFailover(status), "status %d should not trigger failover", status) + } +} From 4731edac6ba54f11e3a9f37d66bb711cc2d8bca1 Mon Sep 17 00:00:00 2001 From: "allen.wu" Date: Tue, 22 Sep 2026 18:42:37 +0800 Subject: [PATCH 2/2] fix(node): harden the L1 RPC failover transport A round of review findings on node/rpcfailover, plus failback to the primary. Endpoint health. 401 and 403 now count as endpoint failures: endpoints carry their own credentials, so an expired key or an exhausted quota on one provider says nothing about the next, and hosted providers use exactly these statuses for it. 3xx counts too, because attempt() rewrites every request's URL: a redirect target would be discarded and the same endpoint asked again until http.Client gives up after ten hops, with no failover. 400 and 404 still pass through, since those describe the request rather than the endpoint. Credential redaction. Two paths leaked an API key held in the URL path or query, because net/http's stripPassword only masks userinfo. One is url.Parse's error, which formats as parse "". The other is the bootstrap URL given to rpc.DialOptions, which is also the URL http.Client reports when it wraps a transport error; that URL is only a default target which attempt() always replaces, so redacting it costs nothing. Response body bound. A server that sends headers and then goes quiet had no bound at all: ResponseHeaderTimeout stops at the headers, the rpc layer sets no deadline, and the node's callers pass contexts derived from context.Background(). The caller blocked forever, and since L1Tracker calls the RPC inline from its tick loop and evaluates the halt gate inside that same loop, the gate stayed open -- the node would keep producing against an arbitrarily stale L1 view with nothing to report it. An inactivity deadline on the connection bounds this without capping legitimately large replies. HTTP/2 is disabled because that bound depends on it. The deadline lives on the connection, and under HTTP/2 one connection carries many concurrent streams whose traffic would keep pushing it forward. TLSNextProto alone is not enough: it stops net/http from handing the connection to its HTTP/2 implementation but not ALPN from negotiating h2, and a server that then selects h2 leaves the HTTP/1.x code reading frames as a response. Failback. A fallback is no longer permanent -- every failbackAfter, one call re-tries the primary, so a recovered primary is re-adopted instead of waiting for the fallback to fail in turn. No background goroutine and no probe traffic; the window is advanced by CompareAndSwap, so concurrent calls share one extra attempt per window rather than paying one each. Each attempt now carries GetBody, so net/http can re-send a request whose pooled connection died. Without it that routine keep-alive race surfaced as an endpoint failure and the sticky switch demoted a healthy primary for the life of the process. Known gaps are unchanged and documented in the package comment: a body that stalls after its first bytes is bounded but cannot be failed over, and a gateway reply that is JSON without being a JSON-RPC envelope is not detected. AI-Developer: claude-opus-4.8 AI-Agent: claude code AI-Reviewer: none Harness-Skill: none Co-Authored-By: Claude Opus 4.8 --- node/rpcfailover/failover.go | 258 ++++++++++++++++--- node/rpcfailover/failover_test.go | 410 ++++++++++++++++++++++++++++-- 2 files changed, 614 insertions(+), 54 deletions(-) diff --git a/node/rpcfailover/failover.go b/node/rpcfailover/failover.go index c74265afb..012687f54 100644 --- a/node/rpcfailover/failover.go +++ b/node/rpcfailover/failover.go @@ -68,6 +68,7 @@ package rpcfailover import ( "bytes" "context" + "crypto/tls" "errors" "fmt" "io" @@ -83,21 +84,69 @@ import ( tmlog "github.com/tendermint/tendermint/libs/log" ) -const ( - // dialTimeout bounds TCP+TLS setup per attempt, so an endpoint whose host is - // blackholed cannot consume the caller's entire deadline. - dialTimeout = 5 * time.Second +// defaultFailbackAfter is how long to keep using a fallback endpoint before +// spending one call to see whether the primary is back. "Primary first" is a +// priority, so a recovered primary should be re-adopted rather than waiting for +// the fallback to fail in turn. +// +// Deliberately not a background health check: no goroutine to manage, no probe +// traffic, and the cost of a still-dead primary is one failed attempt per window +// rather than per call. +// +// Separate from attemptTimeouts because it is failover policy, not a bound on a +// single attempt. +const defaultFailbackAfter = 5 * time.Minute + +// attemptTimeouts bounds one attempt against one endpoint, phase by phase. They +// travel together because they are only meaningful as a set — what is not +// covered by one has to be covered by the next — and because a test that needs +// to shrink one should not have to restate the others. +type attemptTimeouts struct { + // dial bounds DNS resolution and TCP setup, so an endpoint whose host is + // blackholed cannot consume the caller's entire deadline. TLS setup is bounded + // separately by the inherited TLSHandshakeTimeout of 10s. + dial time.Duration - // responseHeaderTimeout bounds the wait for response headers per attempt. - // This is what makes failover work against an endpoint that accepts the - // connection and then never answers — the common shape of a hung RPC - // provider. Because the base transport enforces it per RoundTrip, each - // attempt gets its own budget without wrapping the caller's context. + // responseHeader bounds the wait for response headers. This is what makes + // failover work against an endpoint that accepts the connection and then never + // answers — the common shape of a hung RPC provider. The base transport + // enforces it per RoundTrip, so each attempt gets its own budget without + // wrapping the caller's context. // - // It does not cover a server that sends headers and then stalls the body; - // that case remains bounded only by the caller's context. - responseHeaderTimeout = 10 * time.Second -) + // It does not cover a server that sends headers and then stalls the body — + // net/http documents that it "does not include the time to read the response + // body". idleRead is what bounds that. + responseHeader time.Duration + + // idleRead bounds how long a read may make no progress at all, and is the only + // bound on the response body: responseHeader stops at the headers, the rpc + // layer sets no deadline, and the node's callers pass contexts derived from + // context.Background(). Without it a server that sends headers and then goes + // quiet blocks its caller forever, which matters more than it looks: L1Tracker + // calls the RPC inline from its tick loop and evaluates the halt gate inside + // that same loop, so a call that never returns leaves the gate stuck open — + // the node keeps producing against an arbitrarily stale L1 view with nothing + // to report it. + // + // An inactivity bound, not a total one: a legitimately large reply (a wide + // eth_getLogs can be megabytes) is never cut off while bytes keep arriving. + // + // It also applies to a pooled connection waiting for its next response, so it + // doubles as the idle-connection lifetime. Kept below the inherited + // IdleConnTimeout of 90s so that relationship stays one-way. + // + // The bound is per connection, so it only holds while one request occupies a + // connection at a time. That is why newBaseTransport disables HTTP/2. + idleRead time.Duration +} + +func defaultAttemptTimeouts() attemptTimeouts { + return attemptTimeouts{ + dial: 5 * time.Second, + responseHeader: 10 * time.Second, + idleRead: 60 * time.Second, + } +} // Dial builds a client from a comma-separated list of endpoints, in priority // order with the primary first. @@ -125,12 +174,20 @@ func Dial(ctx context.Context, name, raw string, log tmlog.Logger) (*ethclient.C for _, part := range parts { u, err := url.Parse(part) if err != nil { - return nil, fmt.Errorf("%s: invalid rpc endpoint %s: %w", name, redactEndpoint(part), err) + // url.Error.Error() formats as `parse "": ...`, so wrapping it + // would put the very credentials redactEndpoint exists to hide into the + // log. Report that parsing failed and nothing more. + return nil, fmt.Errorf("%s: rpc endpoint %s is not a valid URL", name, redactEndpoint(part)) } if u.Scheme != "http" && u.Scheme != "https" { return nil, fmt.Errorf("%s: rpc endpoint %s: failover across multiple endpoints requires http(s), got scheme %q", name, redactEndpoint(part), u.Scheme) } + // Hostname(), not Host: "http://:8545" parses with a non-empty Host of + // ":8545" and no host at all. + if u.Hostname() == "" { + return nil, fmt.Errorf("%s: rpc endpoint %s has no host", name, redactEndpoint(part)) + } if u.User != nil { manageAuth = true } @@ -139,14 +196,25 @@ func Dial(ctx context.Context, name, raw string, log tmlog.Logger) (*ethclient.C } transport := &failoverTransport{ - name: name, - base: newBaseTransport(), - endpoints: endpoints, - redacted: redacted, - manageAuth: manageAuth, - log: log, + name: name, + base: newBaseTransport(defaultAttemptTimeouts()), + endpoints: endpoints, + redacted: redacted, + manageAuth: manageAuth, + failbackAfter: defaultFailbackAfter, + log: log, } - rpcClient, err := rpc.DialOptions(ctx, parts[0], rpc.WithHTTPClient(&http.Client{Transport: transport})) + transport.lastProbe.Store(time.Now().UnixNano()) + // The URL handed to DialOptions is only a bootstrap: it selects the transport + // by scheme and becomes the default target, which attempt() always replaces. + // Give it a redacted one, because it is also the URL http.Client reports when + // wrapping a transport error — and stripPassword only masks userinfo, leaving + // an API key in the path or query (Infura /v3/, Alchemy /v2/) to travel into + // every "failed to get L1 header" log line. + // + // This also means http.Client.send no longer stamps basic auth from the + // bootstrap userinfo, which is fine: applyBasicAuth derives it per endpoint. + rpcClient, err := rpc.DialOptions(ctx, redacted[0], rpc.WithHTTPClient(&http.Client{Transport: transport})) if err != nil { return nil, err } @@ -185,15 +253,57 @@ func redactEndpoint(raw string) string { return parsed.Scheme + "://" + parsed.Host } -func newBaseTransport() *http.Transport { +// newBaseTransport builds the transport each attempt runs on. +// +// HTTP/2 is switched off deliberately, and idleRead depends on it. The idle bound +// lives on the TCP connection, and under HTTP/2 one connection carries many +// concurrent streams: traffic on any of them would push the deadline forward, so +// a single stalled response stream would go unbounded again. The node polls L1 +// from several goroutines at once, so that is not hypothetical. Multiplexing buys +// nothing here — these are small, infrequent requests — whereas HTTP/1.1 runs one +// request per connection at a time, which is what makes a per-connection deadline +// mean what it says. A non-nil TLSNextProto is the documented way to disable it. +func newBaseTransport(to attemptTimeouts) *http.Transport { tr := http.DefaultTransport.(*http.Transport).Clone() - tr.DialContext = (&net.Dialer{Timeout: dialTimeout, KeepAlive: 30 * time.Second}).DialContext - tr.ResponseHeaderTimeout = responseHeaderTimeout + tr.ForceAttemptHTTP2 = false + tr.TLSNextProto = map[string]func(string, *tls.Conn) http.RoundTripper{} + // TLSNextProto only stops net/http from handing the connection to its HTTP/2 + // implementation; it does not stop ALPN from negotiating h2 in the first place. + // If the TLS config still advertises h2 the server may select it, and then the + // HTTP/1.x code reads frames as a response and reports the connection broken. + // Pinning NextProtos means h2 can never be selected, whatever else is set here. + tr.TLSClientConfig = &tls.Config{NextProtos: []string{"http/1.1"}} + dialer := &net.Dialer{Timeout: to.dial, KeepAlive: 30 * time.Second} + tr.DialContext = func(ctx context.Context, network, addr string) (net.Conn, error) { + conn, err := dialer.DialContext(ctx, network, addr) + if err != nil { + return nil, err + } + return &idleReadConn{Conn: conn, idle: to.idleRead}, nil + } + tr.ResponseHeaderTimeout = to.responseHeader return tr } -// failoverTransport sends each JSON-RPC request to the endpoint it is currently -// stuck to, and walks the remaining endpoints in ring order when that one fails. +// idleReadConn fails a read that makes no progress for idle. The deadline is +// pushed forward before every read, so it measures inactivity rather than total +// time. +type idleReadConn struct { + net.Conn + idle time.Duration +} + +func (c *idleReadConn) Read(b []byte) (int, error) { + if err := c.Conn.SetReadDeadline(time.Now().Add(c.idle)); err != nil { + return 0, err + } + return c.Conn.Read(b) +} + +// failoverTransport sends each JSON-RPC request to the endpoint currently in +// use, walks the remaining endpoints in ring order when that one fails, and +// every failbackAfter lets one call re-try the primary so that a recovered +// primary is re-adopted rather than waiting for the fallback to fail too. type failoverTransport struct { // name labels logs and errors with which set of endpoints these are. name string @@ -210,13 +320,40 @@ type failoverTransport struct { // cur is the index of the endpoint that last answered. Sticking to it // matters: without it every call would pay the dead primary's timeout again. - // Nothing ever moves it back on its own — a recovered primary is picked up - // when the current endpoint fails, or on the next process restart. cur atomic.Int32 + // failbackAfter is how long to stay on a fallback before letting one call try + // the primary again. A field rather than a constant so tests can shrink it. + failbackAfter time.Duration + + // lastProbe is the unix-nano time the primary was last tried, seeded at + // construction because the process starts out using it. Compared against + // failbackAfter to decide when one more probe is due, and advanced only by + // dueForPrimaryProbe, with CompareAndSwap so exactly one concurrent call pays + // for it. + // + // Switching endpoints deliberately does not touch it. Stamping it on a switch + // would mean fallbacks flapping faster than failbackAfter keep pushing the + // next probe out and the primary is never retried; the cost of not stamping is + // that a switch occurring after the window has already elapsed is followed by + // one immediate re-probe, which also means a primary that only blipped is + // picked back up at once instead of after a full window. + lastProbe atomic.Int64 + log tmlog.Logger } +// dueForPrimaryProbe reports whether enough time has passed to spend one call +// re-trying the primary. Only the caller that wins the CompareAndSwap probes, so +// a fleet of concurrent requests still costs a single extra attempt per window. +func (t *failoverTransport) dueForPrimaryProbe(now int64) bool { + last := t.lastProbe.Load() + if now-last < int64(t.failbackAfter) { + return false + } + return t.lastProbe.CompareAndSwap(last, now) +} + func (t *failoverTransport) RoundTrip(req *http.Request) (*http.Response, error) { // rpc/http.go sets req.Body but leaves req.GetBody nil, so the body must be // buffered here to be replayable against the next endpoint. Request bodies @@ -232,17 +369,41 @@ func (t *failoverTransport) RoundTrip(req *http.Request) (*http.Response, error) } } - start := int(t.cur.Load()) + // prev is the endpoint in use when this call started; start is where this call + // begins walking the ring. They differ only when a failback probe is due, in + // which case this one call starts from the primary instead. If the primary is + // still down the probe costs one failed attempt and the ring carries on. + prev := int(t.cur.Load()) + start := prev + if prev != 0 && t.dueForPrimaryProbe(time.Now().UnixNano()) { + start = 0 + } + var lastErr error for i := range t.endpoints { idx := (start + i) % len(t.endpoints) resp, err := t.attempt(req, body, idx) if err == nil { - if idx != start { + // Compared against prev, not start, so a successful failback probe moves + // back to the primary. + // + // cur is an approximation on purpose — calls do not coordinate, and no + // scheme for agreeing on one endpoint is worth its cost here, because a + // wrong cur only makes the next call spend one failed attempt before + // moving on. Two consequences worth knowing rather than fixing: a call + // that began before someone else switched can write its older index over + // the newer one, and a call whose endpoint equals the one it started from + // records nothing, so a working primary is not re-adopted here — the + // failback probe is what does that. + // + // cause is nil when this was a failback probe, since nothing failed on the + // way here. Recovery is therefore logged at Error with a nil cause; the + // from/to pair is what distinguishes it from a degradation. + if idx != prev { t.cur.Store(int32(idx)) t.log.Error("switched rpc endpoint", "target", t.name, - "from", t.redacted[start], "to", t.redacted[idx], "cause", lastErr) + "from", t.redacted[prev], "to", t.redacted[idx], "cause", lastErr) } return resp, nil } @@ -308,10 +469,10 @@ const peekLimit = 8 // decodes the whole response — full bodies are never buffered, since a single // eth_getLogs reply can be megabytes. func requireJSONBody(resp *http.Response) error { - // Only a success is expected to carry a JSON-RPC reply. The statuses that - // reach here otherwise are the 3xx that http.Client will follow and the 4xx - // that shouldFailover deliberately passes through; neither promises a JSON - // body, and judging them here would undo that decision. + // Only a success is expected to carry a JSON-RPC reply. The only statuses that + // reach here otherwise are the 4xx shouldFailover deliberately passes through, + // which describe the request rather than the endpoint and promise nothing about + // the body; judging them here would undo that decision. if resp.StatusCode < 200 || resp.StatusCode > 299 { return nil } @@ -364,13 +525,32 @@ func replayBody(peeked []byte, rest io.ReadCloser) io.ReadCloser { // shouldFailover reports whether an HTTP status justifies trying the next // endpoint. 5xx and 429 are the endpoint's problem, and 408 means it gave up on -// us. The remaining 4xx are deliberately excluded: a 400/401/403/404 describes -// the request or the credentials, the next endpoint would almost certainly -// reproduce it, and moving on silently would turn a misconfiguration into a -// mystery. A JSON-RPC error arrives as 200 and is not examined here. +// us. +// +// 401 and 403 count too, because endpoints can carry their own credentials (see +// applyBasicAuth): an expired key or an exhausted quota on one provider says +// nothing about the next one, and hosted providers use exactly these statuses +// for it. Refusing to move would blind the node while a correctly configured +// fallback sat idle. A genuine misconfiguration is still visible — every switch +// is logged with its cause, and if all endpoints reject us the aggregate error +// surfaces. +// +// 3xx counts as well. A JSON-RPC endpoint has no business redirecting, and +// following one is not possible here: attempt() rewrites every request's URL to +// its endpoint, so the redirect target would be discarded and the same endpoint +// asked again until http.Client gives up after ten hops — ten wasted round trips +// and no failover. An endpoint configured as http:// behind a server that +// redirects to https:// is the realistic way to hit this. +// +// 400 and 404 are excluded: those describe the request itself, so the next +// endpoint would reproduce them and moving on would only multiply the load. A +// JSON-RPC error arrives as 200 and is not examined here. func shouldFailover(status int) bool { return status == http.StatusRequestTimeout || status == http.StatusTooManyRequests || + status == http.StatusUnauthorized || + status == http.StatusForbidden || + (status >= 300 && status < 400) || status >= 500 } diff --git a/node/rpcfailover/failover_test.go b/node/rpcfailover/failover_test.go index 1354b215c..d89f55482 100644 --- a/node/rpcfailover/failover_test.go +++ b/node/rpcfailover/failover_test.go @@ -23,10 +23,11 @@ const probeBody = `{"jsonrpc":"2.0","id":1,"method":"eth_blockNumber","params":[ // recorder is a test endpoint that counts requests and remembers what it saw. type recorder struct { - server *httptest.Server - hits atomic.Int32 - bodies chan string - auths chan string + server *httptest.Server + hits atomic.Int32 + healthy atomic.Bool // used by newRecoverableEndpoint to bring an endpoint back + bodies chan string + auths chan string } // newEndpoint starts a test endpoint that replies with the given status. A 200 @@ -84,14 +85,17 @@ func newTestTransport(t *testing.T, rawURLs ...string) *failoverTransport { endpoints = append(endpoints, u) redacted = append(redacted, redactEndpoint(raw)) } - return &failoverTransport{ - name: "L1", - base: newBaseTransport(), - endpoints: endpoints, - redacted: redacted, - manageAuth: manageAuth, - log: tmlog.NewNopLogger(), + tr := &failoverTransport{ + name: "L1", + base: newBaseTransport(defaultAttemptTimeouts()), + endpoints: endpoints, + redacted: redacted, + manageAuth: manageAuth, + failbackAfter: defaultFailbackAfter, + log: tmlog.NewNopLogger(), } + tr.lastProbe.Store(time.Now().UnixNano()) // as Dial does + return tr } // post issues a request shaped like the one rpc/http.go builds: a POST whose body @@ -344,14 +348,16 @@ func TestFailoverOnHTMLBlockPageWithoutContentType(t *testing.T) { // to 2xx only, so a 4xx with an empty or HTML body is still returned verbatim // rather than being turned into a failover. func TestClientErrorBodyNotJudged(t *testing.T) { - primary := newRawEndpoint(t, http.StatusForbidden, "text/html", "forbidden") + // 404 rather than 403: 403 now fails over on its own, which would hide what + // this test is about. + primary := newRawEndpoint(t, http.StatusNotFound, "text/html", "not found") secondary := newEndpoint(t, http.StatusOK) tr := newTestTransport(t, primary.url(""), secondary.url("")) resp, err := post(t, tr) require.NoError(t, err) defer resp.Body.Close() - assert.Equal(t, http.StatusForbidden, resp.StatusCode) + assert.Equal(t, http.StatusNotFound, resp.StatusCode) assert.Equal(t, int32(0), secondary.hits.Load()) } @@ -464,6 +470,375 @@ type roundTripFunc func(*http.Request) (*http.Response, error) func (f roundTripFunc) RoundTrip(req *http.Request) (*http.Response, error) { return f(req) } +// --- #4: credentials are per endpoint, so 401/403 must move on --- + +// TestFailoverOnUnauthorized covers an expired or revoked key on the primary. +// Endpoints can carry their own credentials, so 401 on one says nothing about +// the next; refusing to move would blind the caller while a working fallback +// sat idle. +func TestFailoverOnUnauthorized(t *testing.T) { + primary := newEndpoint(t, http.StatusUnauthorized) + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url("alice:expired"), secondary.url("bob:valid")) + + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + + assert.Equal(t, http.StatusOK, resp.StatusCode) + assert.Equal(t, int32(1), secondary.hits.Load()) +} + +// TestFailoverOnForbidden covers an exhausted quota, which hosted providers +// report as 403. +func TestFailoverOnForbidden(t *testing.T) { + primary := newEndpoint(t, http.StatusForbidden) + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + assert.Equal(t, int32(1), secondary.hits.Load()) +} + +// TestNoFailoverOnBadRequestOrNotFound keeps the other half of the rule: these +// describe the request, so another endpoint would only reproduce them. +func TestNoFailoverOnBadRequestOrNotFound(t *testing.T) { + for _, status := range []int{http.StatusBadRequest, http.StatusNotFound} { + primary := newEndpoint(t, status) + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + + resp, err := post(t, tr) + require.NoError(t, err) + resp.Body.Close() + assert.Equal(t, status, resp.StatusCode) + assert.Equal(t, int32(0), secondary.hits.Load(), "status %d must not fail over", status) + } +} + +// --- #5: a malformed endpoint must not leak its credentials into the error --- + +func TestDialErrorDoesNotLeakCredentials(t *testing.T) { + const secret = "s3cr3t-api-key" + // A control character makes url.Parse fail; its own error text embeds the raw + // URL, which is exactly what must not reach the caller. + bad := "http://user:" + secret + "@rpc.example.com/v3/" + secret + "/\x7f" + + _, err := Dial(context.Background(), "L1", bad+",http://ok.invalid", tmlog.NewNopLogger()) + require.Error(t, err) + assert.NotContains(t, err.Error(), secret, "error must not echo credentials from the raw URL") + assert.Contains(t, err.Error(), "is not a valid URL") +} + +func TestDialRejectsEndpointWithoutHost(t *testing.T) { + // "http://:8545" parses with a non-empty Host of ":8545" and no host at all, + // so checking Host alone would let it through. + for _, bad := range []string{"http://", "http://:8545"} { + _, err := Dial(context.Background(), "L1", bad+",http://ok.invalid", tmlog.NewNopLogger()) + require.Error(t, err, "endpoint %q must be rejected", bad) + assert.Contains(t, err.Error(), "has no host") + } +} + +// TestRuntimeErrorDoesNotLeakCredentials covers the error path that outlives +// startup: when every endpoint fails, http.Client wraps the transport error in a +// url.Error carrying the request URL, and its stripPassword only masks userinfo — +// an API key in the path would ride along into every caller's error log. +func TestRuntimeErrorDoesNotLeakCredentials(t *testing.T) { + const secret = "s3cr3t-api-key" + a := newEndpoint(t, http.StatusOK) + b := newEndpoint(t, http.StatusOK) + // The key in the path is how hosted providers embed it (Infura /v3/, Alchemy /v2/). + urlA := a.url("") + "/v3/" + secret + urlB := b.url("") + "/v3/" + secret + a.server.Close() // both refuse connections, so the call fails at runtime + b.server.Close() + + client, err := Dial(context.Background(), "L1", urlA+","+urlB, tmlog.NewNopLogger()) + require.NoError(t, err) + defer client.Close() + + _, err = client.BlockNumber(context.Background()) + require.Error(t, err) + assert.NotContains(t, err.Error(), secret, + "the URL http.Client reports on a transport error must not carry the API key") +} + +// TestFailoverOnRedirect: attempt() rewrites every request's URL to its endpoint, +// so a redirect target would be discarded and the same endpoint asked again until +// http.Client gives up after ten hops. Treating 3xx as an endpoint failure is +// what keeps that from happening. +func TestFailoverOnRedirect(t *testing.T) { + redirecting := &recorder{bodies: make(chan string, 8), auths: make(chan string, 8)} + redirecting.server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + redirecting.hits.Add(1) + w.Header().Set("Location", "https://elsewhere.invalid/rpc") + w.WriteHeader(http.StatusFound) + })) + t.Cleanup(redirecting.server.Close) + + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, redirecting.url(""), secondary.url("")) + + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + + assert.Equal(t, http.StatusOK, resp.StatusCode) + assert.Equal(t, int32(1), redirecting.hits.Load(), "the redirect must be asked once, not looped over") + assert.Equal(t, int32(1), secondary.hits.Load()) +} + +// --- #6: a recovered primary is re-adopted --- + +func TestFailsBackToPrimaryAfterWindow(t *testing.T) { + primary := newRecoverableEndpoint(t) + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + tr.failbackAfter = time.Millisecond + + primary.healthy.Store(false) + resp, err := post(t, tr) + require.NoError(t, err) + resp.Body.Close() + require.Equal(t, int32(1), tr.cur.Load(), "should have moved to the secondary") + + primary.healthy.Store(true) + time.Sleep(2 * time.Millisecond) // let the failback window elapse + + resp, err = post(t, tr) + require.NoError(t, err) + resp.Body.Close() + + assert.Equal(t, int32(0), tr.cur.Load(), "a healthy primary must be re-adopted") + assert.Equal(t, int32(1), secondary.hits.Load(), "the secondary should not have been needed again") +} + +// TestNormalCallDoesNotTouchProbeWindow pins that the failback clock is only +// advanced when the primary is actually tried: if every successful call restamped +// it, a busy node would never reach the window and the primary would never be +// retried. +func TestNormalCallDoesNotTouchProbeWindow(t *testing.T) { + only := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, only.url(""), only.url("")) + tr.lastProbe.Store(12345) + + for i := 0; i < 3; i++ { + resp, err := post(t, tr) + require.NoError(t, err) + resp.Body.Close() + } + + assert.Equal(t, int64(12345), tr.lastProbe.Load(), + "a call served by the endpoint already in use must not move the failback clock") +} + +// TestSwitchDoesNotPostponeProbe: only dueForPrimaryProbe advances the failback +// clock. If switching endpoints stamped it too, fallbacks flapping faster than +// failbackAfter would keep pushing the next probe out and the primary would never +// be retried. +func TestSwitchDoesNotPostponeProbe(t *testing.T) { + primary := newEndpoint(t, http.StatusServiceUnavailable) + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + + // Long window, so no probe is due during this call and the clock can only + // change if the switch itself stamps it. + tr.failbackAfter = time.Hour + stamp := time.Now().Add(-time.Minute).UnixNano() + tr.lastProbe.Store(stamp) + + resp, err := post(t, tr) + require.NoError(t, err) + resp.Body.Close() + + require.Equal(t, int32(1), tr.cur.Load(), "should have switched to the secondary") + assert.Equal(t, stamp, tr.lastProbe.Load(), "a switch must leave the failback clock alone") +} + +// TestSwitchAfterWindowElapsedReprobesOnce is the accepted cost of the above: a +// switch that happens once the window has already passed is followed by one +// immediate re-probe. That also means a primary which only blipped is picked back +// up at once rather than after a full window. +func TestSwitchAfterWindowElapsedReprobesOnce(t *testing.T) { + primary := newRecoverableEndpoint(t) + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + tr.failbackAfter = time.Millisecond + + primary.healthy.Store(false) + resp, err := post(t, tr) + require.NoError(t, err) + resp.Body.Close() + require.Equal(t, int32(1), tr.cur.Load()) + require.Equal(t, int32(1), primary.hits.Load()) + + time.Sleep(2 * time.Millisecond) // window elapses while still on the fallback + + // Primary came back between the two calls, so the re-probe adopts it again. + primary.healthy.Store(true) + resp, err = post(t, tr) + require.NoError(t, err) + resp.Body.Close() + + assert.Equal(t, int32(2), primary.hits.Load(), "the primary is re-probed once the window has passed") + assert.Equal(t, int32(0), tr.cur.Load(), "and re-adopted now that it answers") +} + +func TestDoesNotProbePrimaryEveryCallWhileItIsDown(t *testing.T) { + primary := newEndpoint(t, http.StatusServiceUnavailable) + secondary := newEndpoint(t, http.StatusOK) + tr := newTestTransport(t, primary.url(""), secondary.url("")) + // Long window: the first call switches away, later calls must not re-probe. + tr.failbackAfter = time.Hour + + for i := 0; i < 4; i++ { + resp, err := post(t, tr) + require.NoError(t, err) + resp.Body.Close() + } + + assert.Equal(t, int32(1), primary.hits.Load(), + "the dead primary is probed once, not on every call") + assert.Equal(t, int32(4), secondary.hits.Load()) +} + +// newRecoverableEndpoint serves 503 or a valid reply depending on a flag, so a +// test can bring an endpoint back up. +func newRecoverableEndpoint(t *testing.T) *recorder { + t.Helper() + r := &recorder{bodies: make(chan string, 8), auths: make(chan string, 8)} + r.healthy.Store(true) + r.server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + r.hits.Add(1) + if !r.healthy.Load() { + w.WriteHeader(http.StatusServiceUnavailable) + return + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"jsonrpc":"2.0","id":1,"result":"0x10"}`)) + })) + t.Cleanup(r.server.Close) + return r +} + +// --- a body that stops arriving must be bounded, not block forever --- + +// newStallingEndpoint starts an endpoint that writes its headers and prefix, +// flushes them onto the wire, and then stops writing without closing the +// connection — the shape neither ResponseHeaderTimeout nor a TCP error catches. +// The returned channel must be closed by the test before its server is torn +// down, since httptest.Server.Close waits for outstanding handlers. +func newStallingEndpoint(t *testing.T, prefix string) (*recorder, chan struct{}) { + t.Helper() + release := make(chan struct{}) + r := &recorder{bodies: make(chan string, 8), auths: make(chan string, 8)} + r.server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + r.hits.Add(1) + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + if prefix != "" { + _, _ = io.WriteString(w, prefix) + } + w.(http.Flusher).Flush() + <-release + })) + t.Cleanup(r.server.Close) + return r, release +} + +// TestFailoverWhenBodyNeverArrives: the stall lands inside requireJSONBody's +// peek, which is still inside the attempt, so it is a normal endpoint failure and +// the next endpoint serves the call. +func TestFailoverWhenBodyNeverArrives(t *testing.T) { + primary, release := newStallingEndpoint(t, "") + defer close(release) + secondary := newEndpoint(t, http.StatusOK) + + tr := newTestTransport(t, primary.url(""), secondary.url("")) + to := defaultAttemptTimeouts() + to.idleRead = 100 * time.Millisecond + tr.base = newBaseTransport(to) + + started := time.Now() + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + + assert.Equal(t, int32(1), primary.hits.Load()) + assert.Equal(t, int32(1), secondary.hits.Load(), "a body that never arrives is an endpoint failure") + assert.Less(t, time.Since(started), 5*time.Second, "must not have waited indefinitely on the stalled endpoint") + + body, err := io.ReadAll(resp.Body) + require.NoError(t, err) + assert.Contains(t, string(body), `"result":"0x10"`) +} + +// TestStalledBodyMidStreamFailsBounded is the residual case: enough bytes arrive +// to satisfy the peek, so the response is handed up and the stall happens past +// the failover decision point. No switch is possible there — but the read must +// fail rather than hang, which is what lets L1Tracker's loop iterate again and +// its halt gate keep working. +func TestStalledBodyMidStreamFailsBounded(t *testing.T) { + primary, release := newStallingEndpoint(t, `{"jsonrpc":"2.0",`) + defer close(release) + secondary := newEndpoint(t, http.StatusOK) + + tr := newTestTransport(t, primary.url(""), secondary.url("")) + to := defaultAttemptTimeouts() + to.idleRead = 100 * time.Millisecond + tr.base = newBaseTransport(to) + + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + require.Equal(t, int32(0), secondary.hits.Load(), "the peek succeeded, so this call does not fail over") + + started := time.Now() + _, err = io.ReadAll(resp.Body) + require.Error(t, err, "a stalled body must fail rather than block forever") + assert.Less(t, time.Since(started), 5*time.Second) +} + +// TestHTTP2IsDisabled pins the assumption idleRead rests on. Under HTTP/2 one +// connection carries many concurrent streams, so traffic on any of them would +// push a per-connection read deadline forward and a single stalled stream would +// be unbounded again. The server here offers h2 over ALPN; the transport must +// still come back on HTTP/1.1. +func TestHTTP2IsDisabled(t *testing.T) { + srv := httptest.NewUnstartedServer(http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"jsonrpc":"2.0","id":1,"result":"0x10"}`)) + })) + srv.EnableHTTP2 = true + srv.StartTLS() + t.Cleanup(srv.Close) + + // Control: without this the assertion below would also pass against a server + // that never offered h2 in the first place. + ctrl, err := srv.Client().Post(srv.URL, "application/json", strings.NewReader(probeBody)) + require.NoError(t, err) + require.NoError(t, ctrl.Body.Close()) + require.Equal(t, "HTTP/2.0", ctrl.Proto, "precondition: the test server must offer HTTP/2") + + tr := newTestTransport(t, srv.URL, srv.URL) + base := newBaseTransport(defaultAttemptTimeouts()) + // Trust the test server's certificate, but only that: replacing the whole TLS + // config would re-advertise h2 over ALPN, and the server would then select a + // protocol the HTTP/1.x code cannot read. + base.TLSClientConfig.RootCAs = srv.Client().Transport.(*http.Transport).TLSClientConfig.RootCAs + tr.base = base + + resp, err := post(t, tr) + require.NoError(t, err) + defer resp.Body.Close() + + assert.Equal(t, "HTTP/1.1", resp.Proto, + "HTTP/2 must stay disabled: idleRead is a per-connection bound and needs one request per connection") +} + func TestFailoverOnTooManyRequests(t *testing.T) { primary := newEndpoint(t, http.StatusTooManyRequests) secondary := newEndpoint(t, http.StatusOK) @@ -619,10 +994,15 @@ func TestRedactedEndpointsAreDistinct(t *testing.T) { } func TestShouldFailover(t *testing.T) { - for _, status := range []int{408, 429, 500, 502, 503, 504} { + // 401 and 403 are in this list because endpoints can carry their own + // credentials: an expired key or exhausted quota on one provider says nothing + // about the next. 3xx is here because attempt() rewrites the URL, so a redirect + // could never be followed — only looped over. + for _, status := range []int{408, 429, 401, 403, 301, 302, 307, 308, 500, 502, 503, 504} { assert.True(t, shouldFailover(status), "status %d should trigger failover", status) } - for _, status := range []int{200, 201, 301, 400, 401, 403, 404} { + // These describe the request, so another endpoint would only reproduce them. + for _, status := range []int{200, 201, 400, 404} { assert.False(t, shouldFailover(status), "status %d should not trigger failover", status) } }