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..012687f54 --- /dev/null +++ b/node/rpcfailover/failover.go @@ -0,0 +1,568 @@ +// 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" + "crypto/tls" + "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" +) + +// 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 + + // 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 — + // 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. +// +// 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 { + // 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 + } + endpoints = append(endpoints, u) + redacted = append(redacted, redactEndpoint(part)) + } + + transport := &failoverTransport{ + name: name, + base: newBaseTransport(defaultAttemptTimeouts()), + endpoints: endpoints, + redacted: redacted, + manageAuth: manageAuth, + failbackAfter: defaultFailbackAfter, + log: log, + } + 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 + } + 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 +} + +// 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.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 +} + +// 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 + + 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. + 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 + // 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) + } + } + + // 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 { + // 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[prev], "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 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 + } + + // 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. +// +// 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 +} + +// 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..d89f55482 --- /dev/null +++ b/node/rpcfailover/failover_test.go @@ -0,0 +1,1008 @@ +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 + 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 +// 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)) + } + 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 +// 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) { + // 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.StatusNotFound, 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) } + +// --- #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) + 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) { + // 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) + } + // 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) + } +}