From d4ba577b03a7097e19b01bbbbf4e77a3e7573f27 Mon Sep 17 00:00:00 2001 From: jonny rhea Date: Mon, 29 Dec 2025 12:20:59 -0600 Subject: [PATCH] address review feedback --- cmd/keeper/go.mod | 2 +- cmd/keeper/go.sum | 12 +++---- rpc/client.go | 9 +----- rpc/client_opt.go | 4 --- rpc/handler.go | 23 +++++++++---- rpc/server.go | 27 ++++++---------- rpc/tracing_test.go | 78 +++++++-------------------------------------- 7 files changed, 46 insertions(+), 109 deletions(-) diff --git a/cmd/keeper/go.mod b/cmd/keeper/go.mod index a42be042aa..cee1ce05a7 100644 --- a/cmd/keeper/go.mod +++ b/cmd/keeper/go.mod @@ -34,7 +34,7 @@ require ( github.com/tklauser/numcpus v0.6.1 // indirect golang.org/x/crypto v0.36.0 // indirect golang.org/x/sync v0.12.0 // indirect - golang.org/x/sys v0.36.0 // indirect + golang.org/x/sys v0.39.0 // indirect gopkg.in/yaml.v2 v2.4.0 // indirect ) diff --git a/cmd/keeper/go.sum b/cmd/keeper/go.sum index 133a3b10b1..b93969cc60 100644 --- a/cmd/keeper/go.sum +++ b/cmd/keeper/go.sum @@ -96,12 +96,12 @@ github.com/prometheus/procfs v0.9.0 h1:wzCHvIvM5SxWqYvwgVL7yJY8Lz3PKn49KQtpgMYJf github.com/prometheus/procfs v0.9.0/go.mod h1:+pB4zwohETzFnmlpe6yd2lSc+0/46IYZRB/chUwxUZY= github.com/prysmaticlabs/gohashtree v0.0.4-beta h1:H/EbCuXPeTV3lpKeXGPpEV9gsUpkqOOVnWapUyeWro4= github.com/prysmaticlabs/gohashtree v0.0.4-beta/go.mod h1:BFdtALS+Ffhg3lGQIHv9HDWuHS8cTvHZzrHWxwOtGOs= -github.com/rogpeppe/go-internal v1.12.0 h1:exVL4IDcn6na9z1rAb56Vxr+CgyK3nn3O+epU5NdKM8= -github.com/rogpeppe/go-internal v1.12.0/go.mod h1:E+RYuTGaKKdloAfM02xzb0FW3Paa99yedzYV+kq4uf4= +github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= +github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= github.com/shirou/gopsutil v3.21.4-0.20210419000835-c7a38de76ee5+incompatible h1:Bn1aCHHRnjv4Bl16T8rcaFjYSrGrIZvpiGO6P3Q4GpU= github.com/shirou/gopsutil v3.21.4-0.20210419000835-c7a38de76ee5+incompatible/go.mod h1:5b4v6he4MtMOwMlS0TUMTu2PcXUg8+E1lC7eC3UO/RA= -github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA= -github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= github.com/supranational/blst v0.3.16-0.20250831170142-f48500c1fdbe h1:nbdqkIGOGfUAD54q1s2YBcBz/WcsxCO9HUQ4aGV5hUw= github.com/supranational/blst v0.3.16-0.20250831170142-f48500c1fdbe/go.mod h1:jZJtfjgudtNl4en1tzwPIV3KjUnQUvG3/j+w+fVonLw= github.com/tklauser/go-sysconf v0.3.12 h1:0QaGUFOdQaIVdPgfITYzaTegZvdCjmYO52cSFAEVmqU= @@ -118,8 +118,8 @@ golang.org/x/sys v0.0.0-20190916202348-b4ddaad3f8a3/go.mod h1:h1NjWce9XRLGQEsW7w golang.org/x/sys v0.1.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.8.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.11.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.36.0 h1:KVRy2GtZBrk1cBYA7MKu5bEZFxQk4NIDV6RLVcC8o0k= -golang.org/x/sys v0.36.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= +golang.org/x/sys v0.39.0 h1:CvCKL8MeisomCi6qNZ+wbb0DN9E5AATixKsvNtMoMFk= +golang.org/x/sys v0.39.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= golang.org/x/text v0.23.0 h1:D71I7dUrlY+VX0gQShAThNGHFxZ13dGLBHQLVl1mJlY= golang.org/x/text v0.23.0/go.mod h1:/BLNzu4aZCJ1+kcD0DNRotWKage4q2rGVAg4o22unh4= google.golang.org/protobuf v1.34.2 h1:6xV6lTsCfpGD21XK49h7MhtcApnLqkfYgPcdHftf6hg= diff --git a/rpc/client.go b/rpc/client.go index 0d54b9af1b..8d81503d59 100644 --- a/rpc/client.go +++ b/rpc/client.go @@ -29,8 +29,6 @@ import ( "time" "github.com/ethereum/go-ethereum/log" - "go.opentelemetry.io/otel" - "go.opentelemetry.io/otel/trace" ) var ( @@ -90,7 +88,6 @@ type Client struct { // config fields batchItemLimit int batchResponseMaxSize int - tracer trace.Tracer // writeConn is used for writing to the connection on the caller's goroutine. It should // only be accessed outside of dispatch, with the write lock held. The write lock is @@ -122,7 +119,7 @@ func (c *Client) newClientConn(conn ServerCodec) *clientConn { ctx := context.Background() ctx = context.WithValue(ctx, clientContextKey{}, c) ctx = context.WithValue(ctx, peerInfoContextKey{}, conn.peerInfo()) - handler := newHandler(ctx, conn, c.idgen, c.services, c.batchItemLimit, c.batchResponseMaxSize, c.tracer) + handler := newHandler(ctx, conn, c.idgen, c.services, c.batchItemLimit, c.batchResponseMaxSize, nil) return &clientConn{conn, handler} } @@ -250,7 +247,6 @@ func initClient(conn ServerCodec, services *serviceRegistry, cfg *clientConfig) idgen: cfg.idgen, batchItemLimit: cfg.batchItemLimit, batchResponseMaxSize: cfg.batchResponseLimit, - tracer: cfg.tracer, writeConn: conn, close: make(chan struct{}), closing: make(chan struct{}), @@ -267,9 +263,6 @@ func initClient(conn ServerCodec, services *serviceRegistry, cfg *clientConfig) if c.idgen == nil { c.idgen = randomIDGenerator() } - if c.tracer == nil { - c.tracer = otel.GetTracerProvider().Tracer(tracerName) - } // Launch the main loop. if !isHTTP { diff --git a/rpc/client_opt.go b/rpc/client_opt.go index 38ca2d8946..3fa045a9b9 100644 --- a/rpc/client_opt.go +++ b/rpc/client_opt.go @@ -20,7 +20,6 @@ import ( "net/http" "github.com/gorilla/websocket" - "go.opentelemetry.io/otel/trace" ) // ClientOption is a configuration option for the RPC client. @@ -42,9 +41,6 @@ type clientConfig struct { idgen func() ID batchItemLimit int batchResponseLimit int - - // Tracing - tracer trace.Tracer } func (cfg *clientConfig) initHeaders() { diff --git a/rpc/handler.go b/rpc/handler.go index 46eddc1b7c..65f5304480 100644 --- a/rpc/handler.go +++ b/rpc/handler.go @@ -28,12 +28,12 @@ import ( "sync" "time" + "github.com/ethereum/go-ethereum/log" + "go.opentelemetry.io/otel" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/codes" semconv "go.opentelemetry.io/otel/semconv/v1.24.0" "go.opentelemetry.io/otel/trace" - - "github.com/ethereum/go-ethereum/log" ) // handler handles JSON-RPC messages. There is one handler per connection. Note that @@ -70,7 +70,7 @@ type handler struct { allowSubscribe bool batchRequestLimit int batchResponseMaxSize int - tracer trace.Tracer + tracerProvider trace.TracerProvider subLock sync.Mutex serverSubs map[ID]*Subscription @@ -82,7 +82,7 @@ type callProc struct { isBatch bool } -func newHandler(connCtx context.Context, conn jsonWriter, idgen func() ID, reg *serviceRegistry, batchRequestLimit, batchResponseMaxSize int, tracer trace.Tracer) *handler { +func newHandler(connCtx context.Context, conn jsonWriter, idgen func() ID, reg *serviceRegistry, batchRequestLimit, batchResponseMaxSize int, tracerProvider trace.TracerProvider) *handler { rootCtx, cancelRoot := context.WithCancel(connCtx) h := &handler{ reg: reg, @@ -97,7 +97,7 @@ func newHandler(connCtx context.Context, conn jsonWriter, idgen func() ID, reg * log: log.Root(), batchRequestLimit: batchRequestLimit, batchResponseMaxSize: batchResponseMaxSize, - tracer: tracer, + tracerProvider: tracerProvider, } if conn.remoteAddr() != "" { h.log = h.log.New("conn", conn.remoteAddr()) @@ -608,7 +608,7 @@ func (h *handler) handleSubscribe(cp *callProc, msg *jsonrpcMessage) *jsonrpcMes // startRPCSpan starts a tracing span for an RPC call. func (h *handler) startRPCSpan(ctx context.Context, msg *jsonrpcMessage, isBatch bool) (context.Context, trace.Span) { - ctx, span := h.tracer.Start(ctx, "rpc.call") + ctx, span := h.tracer().Start(ctx, "rpc.call") span.SetAttributes( semconv.RPCSystemKey.String("jsonrpc"), semconv.RPCMethodKey.String(msg.Method), @@ -622,6 +622,17 @@ func (h *handler) startRPCSpan(ctx context.Context, msg *jsonrpcMessage, isBatch return ctx, span } +// tracer returns the OpenTelemetry Tracer for RPC call tracing. +func (h *handler) tracer() trace.Tracer { + if h.tracerProvider == nil { + // Default to global TracerProvider if none is set. + // Note: If no TracerProvider is set, the default is a no-op TracerProvider. + // See https://pkg.go.dev/go.opentelemetry.io/otel#GetTracerProvider + return otel.Tracer("") + } + return h.tracerProvider.Tracer("") +} + // runMethod runs the Go callback for an RPC method. func (h *handler) runMethod(ctx context.Context, msg *jsonrpcMessage, callb *callback, args []reflect.Value) *jsonrpcMessage { result, err := callb.call(ctx, msg.Method, args) diff --git a/rpc/server.go b/rpc/server.go index d4c9d1ba4e..ddf396b422 100644 --- a/rpc/server.go +++ b/rpc/server.go @@ -24,15 +24,12 @@ import ( "sync" "sync/atomic" - "go.opentelemetry.io/otel" - "go.opentelemetry.io/otel/trace" - "github.com/ethereum/go-ethereum/log" + "go.opentelemetry.io/otel/trace" ) const MetadataApi = "rpc" const EngineApi = "engine" -const tracerName = "github.com/ethereum/go-ethereum/rpc" // CodecOption specifies which type of messages a codec supports. // @@ -51,7 +48,6 @@ const ( type Server struct { services serviceRegistry idgen func() ID - tracer trace.Tracer mutex sync.Mutex codecs map[ServerCodec]struct{} @@ -60,16 +56,17 @@ type Server struct { batchResponseLimit int httpBodyLimit int wsReadLimit int64 + tracerProvider trace.TracerProvider } // NewServer creates a new server instance with no registered handlers. func NewServer() *Server { server := &Server{ - idgen: randomIDGenerator(), - codecs: make(map[ServerCodec]struct{}), - httpBodyLimit: defaultBodyLimit, - wsReadLimit: wsDefaultReadLimit, - tracer: otel.GetTracerProvider().Tracer(tracerName), + idgen: randomIDGenerator(), + codecs: make(map[ServerCodec]struct{}), + httpBodyLimit: defaultBodyLimit, + wsReadLimit: wsDefaultReadLimit, + tracerProvider: nil, } server.run.Store(true) // Register the default service providing meta information about the RPC service such @@ -106,12 +103,7 @@ func (s *Server) SetWebsocketReadLimit(limit int64) { // SetTracerProvider configures the OpenTelemetry TracerProvider for RPC call tracing. func (s *Server) SetTracerProvider(tp trace.TracerProvider) { - if tp == nil { - // Reset to the global provider (which is noop unless configured elsewhere). - s.tracer = otel.GetTracerProvider().Tracer(tracerName) - return - } - s.tracer = tp.Tracer(tracerName) + s.tracerProvider = tp } // RegisterName creates a service for the given receiver type under the given name. When no @@ -139,7 +131,6 @@ func (s *Server) ServeCodec(codec ServerCodec, options CodecOption) { idgen: s.idgen, batchItemLimit: s.batchItemLimit, batchResponseLimit: s.batchResponseLimit, - tracer: s.tracer, } c := initClient(codec, &s.services, cfg) <-codec.closed() @@ -173,7 +164,7 @@ func (s *Server) serveSingleRequest(ctx context.Context, codec ServerCodec) { return } - h := newHandler(ctx, codec, s.idgen, &s.services, s.batchItemLimit, s.batchResponseLimit, s.tracer) + h := newHandler(ctx, codec, s.idgen, &s.services, s.batchItemLimit, s.batchResponseLimit, s.tracerProvider) h.allowSubscribe = false defer h.close(io.EOF, nil) diff --git a/rpc/tracing_test.go b/rpc/tracing_test.go index 570bb0488b..a1ab9ddab8 100644 --- a/rpc/tracing_test.go +++ b/rpc/tracing_test.go @@ -26,14 +26,6 @@ import ( "go.opentelemetry.io/otel/sdk/trace/tracetest" ) -func spanNames(spans []tracetest.SpanStub) []string { - names := make([]string, len(spans)) - for i, s := range spans { - names[i] = s.Name - } - return names -} - func attributeMap(attrs []attribute.KeyValue) map[string]string { m := make(map[string]string) for _, a := range attrs { @@ -57,14 +49,14 @@ func newTracingServer(t *testing.T) (*Server, *sdktrace.TracerProvider, *tracete t.Helper() exporter := tracetest.NewInMemoryExporter() - tracer := sdktrace.NewTracerProvider(sdktrace.WithSyncer(exporter)) - t.Cleanup(func() { _ = tracer.Shutdown(context.Background()) }) + tp := sdktrace.NewTracerProvider(sdktrace.WithSyncer(exporter)) + t.Cleanup(func() { _ = tp.Shutdown(context.Background()) }) server := newTestServer() - server.SetTracerProvider(tracer) + server.SetTracerProvider(tp) t.Cleanup(server.Stop) - return server, tracer, exporter + return server, tp, exporter } // TestTracingHTTP verifies that RPC spans are emitted when processing HTTP requests. @@ -103,52 +95,7 @@ func TestTracingHTTP(t *testing.T) { } } if rpcSpan == nil { - t.Fatalf("rpc.call span not found, got %v", spanNames(spans)) - } - - attrs := attributeMap(rpcSpan.Attributes) - if attrs["rpc.system"] != "jsonrpc" { - t.Errorf("expected rpc.system=jsonrpc, got %v", attrs["rpc.system"]) - } - if attrs["rpc.method"] != "test_echo" { - t.Errorf("expected rpc.method=test_echo, got %v", attrs["rpc.method"]) - } - if _, ok := attrs["rpc.id"]; !ok { - t.Errorf("expected rpc.id attribute to be set") - } -} - -// TestTracingInProcess verifies that RPC spans are emitted when processing requests in process. -func TestTracingInProcess(t *testing.T) { - t.Parallel() - server, tracer, exporter := newTracingServer(t) - - client := DialInProc(server) - t.Cleanup(client.Close) - - var result echoResult - if err := client.Call(&result, "test_echo", "hello", 42, &echoArgs{S: "world"}); err != nil { - t.Fatalf("RPC call failed: %v", err) - } - - if err := tracer.ForceFlush(context.Background()); err != nil { - t.Fatalf("failed to flush: %v", err) - } - - spans := exporter.GetSpans() - if len(spans) == 0 { - t.Fatal("no spans were emitted") - } - - var rpcSpan *tracetest.SpanStub - for i := range spans { - if spans[i].Name == "rpc.call" { - rpcSpan = &spans[i] - break - } - } - if rpcSpan == nil { - t.Fatalf("rpc.call span not found, got %v", spanNames(spans)) + t.Fatalf("rpc.call span not found.") } attrs := attributeMap(rpcSpan.Attributes) @@ -204,14 +151,13 @@ func TestTracingBatchHTTP(t *testing.T) { var found int for i := range spans { - if spans[i].Name != "rpc.call" { - continue - } - attrs := attributeMap(spans[i].Attributes) - if attrs["rpc.system"] == "jsonrpc" && - attrs["rpc.method"] == "test_echo" && - attrs["rpc.batch"] == "true" { - found++ + if spans[i].Name == "rpc.call" { + attrs := attributeMap(spans[i].Attributes) + if attrs["rpc.system"] == "jsonrpc" && + attrs["rpc.method"] == "test_echo" && + attrs["rpc.batch"] == "true" { + found++ + } } } if found != len(batch) {