From a7b9db17221c5d2ea56ca34f86012b452fc616b8 Mon Sep 17 00:00:00 2001 From: snowkide Date: Fri, 2 Oct 2026 11:50:25 +0200 Subject: [PATCH 1/2] feat: port CLL auth and jsonrpc passthrough onto rewritten feat/websocket-support Re-applies the net change of #24 (#20 + #21 + #23) onto the 2026-09-28 rewrite of feat/websocket-support (stacked on #29): - auth: forwardedClientId strategy; secret accepted as path /, ?apikey= query and header (NewPayloadFromHttp takes the URL path) - architecture: jsonrpc passthrough for non-EVM HAProxy upstreams, now coexisting with upstream's svm architecture and registered as a no-op ArchitectureHandler (required by the new architecture registry); its error extractor defers to the EVM one, as before - json-rpc: treat "error": null as success; omit empty params Co-Authored-By: Claude Opus 5.5 --- auth/authorizer.go | 5 ++ auth/http.go | 58 +++++++++++- auth/payload.go | 17 ++-- auth/strategy_forwarded_client_id.go | 40 +++++++++ auth/strategy_forwarded_client_id_test.go | 104 ++++++++++++++++++++++ auth/strategy_secret.go | 9 ++ clients/registry.go | 6 +- common/architecture_jsonrpc.go | 56 ++++++++++++ common/architecture_jsonrpc_test.go | 63 +++++++++++++ common/config.go | 49 +++++++--- common/defaults.go | 42 ++++++++- common/json_rpc.go | 24 +++-- common/json_rpc_test.go | 29 +++++- common/network.go | 11 ++- common/request.go | 25 ++++++ common/validation.go | 41 ++++++++- erpc/healthcheck.go | 2 +- erpc/http_server.go | 4 +- erpc/http_server_test.go | 2 +- erpc/networks.go | 6 +- erpc/networks_registry.go | 7 ++ erpc/ws_server.go | 2 +- typescript/config/src/generated.ts | 12 +++ typescript/config/src/index.ts | 1 + upstream/registry.go | 2 + upstream/upstream.go | 8 ++ util/ids.go | 8 ++ 27 files changed, 591 insertions(+), 42 deletions(-) create mode 100644 auth/strategy_forwarded_client_id.go create mode 100644 auth/strategy_forwarded_client_id_test.go create mode 100644 common/architecture_jsonrpc.go create mode 100644 common/architecture_jsonrpc_test.go diff --git a/auth/authorizer.go b/auth/authorizer.go index 464d2b306..0628881e5 100644 --- a/auth/authorizer.go +++ b/auth/authorizer.go @@ -63,6 +63,11 @@ func NewAuthorizer(appCtx context.Context, logger *zerolog.Logger, projectId str if err != nil { return nil, err } + case common.AuthTypeForwardedClientId: + if cfg.ForwardedClientId == nil { + return nil, common.NewErrInvalidConfig("forwardedClientId strategy config is nil") + } + strategy = NewForwardedClientIdStrategy(cfg.ForwardedClientId) default: return nil, common.NewErrInvalidConfig(fmt.Sprintf("unknown auth strategy type: %s", cfg.Type)) } diff --git a/auth/http.go b/auth/http.go index 3b6698662..2ace053e6 100644 --- a/auth/http.go +++ b/auth/http.go @@ -5,12 +5,13 @@ import ( "errors" "net/http" "net/url" + "path" "strings" "github.com/erpc/erpc/common" ) -func NewPayloadFromHttp(method string, remoteAddr string, headers http.Header, args url.Values) (*AuthPayload, error) { +func NewPayloadFromHttp(method string, remoteAddr string, headers http.Header, args url.Values, requestPath string) (*AuthPayload, error) { ap := &AuthPayload{ Method: method, } @@ -25,11 +26,22 @@ func NewPayloadFromHttp(method string, remoteAddr string, headers http.Header, a ap.Secret = &SecretPayload{ Value: secret, } + } else if apikey := args.Get("apikey"); apikey != "" { + // Alias used by edge gateways / clients that speak "apikey" rather than "secret". + ap.Type = common.AuthTypeSecret + ap.Secret = &SecretPayload{ + Value: apikey, + } } else if tkn := headers.Get("X-ERPC-Secret-Token"); tkn != "" { ap.Type = common.AuthTypeSecret ap.Secret = &SecretPayload{ Value: tkn, } + } else if apikey := firstNonEmptyHeader(headers, "apikey", "X-Api-Key"); apikey != "" { + ap.Type = common.AuthTypeSecret + ap.Secret = &SecretPayload{ + Value: apikey, + } } else if ath := headers.Get("Authorization"); ath != "" { ath = strings.TrimSpace(ath) parts := strings.SplitN(ath, " ", 2) @@ -77,6 +89,19 @@ func NewPayloadFromHttp(method string, remoteAddr string, headers http.Header, a Message: normalizeSiweMessage(msg), } } + } else if pathSecret := singlePathSegmentSecret(requestPath); pathSecret != "" { + // Path form: https://host/ (with domain aliasing so the segment + // is not consumed as project/network). Avoids edge Lua/WASM filters. + ap.Type = common.AuthTypeSecret + ap.Secret = &SecretPayload{ + Value: pathSecret, + } + } else if clientId := firstNonEmptyHeader(headers, "X-Client-Id", "x-client-id"); clientId != "" { + // Gateway-injected identity after edge API-key auth (Envoy forwardClientIDHeader). + ap.Type = common.AuthTypeForwardedClientId + ap.ForwardedClientId = &ForwardedClientIdPayload{ + Value: clientId, + } } // Default to network strategy when no other auth signals are present. @@ -87,6 +112,37 @@ func NewPayloadFromHttp(method string, remoteAddr string, headers http.Header, a return ap, nil } +// singlePathSegmentSecret returns the sole path segment when the URL is +// `/` (or `//`). Multi-segment eRPC paths and reserved +// endpoints are ignored so routing/healthchecks stay unchanged. +func singlePathSegmentSecret(requestPath string) string { + if requestPath == "" { + return "" + } + ps := path.Clean(requestPath) + if ps == "/" || ps == "." { + return "" + } + seg := strings.TrimPrefix(ps, "/") + if seg == "" || strings.Contains(seg, "/") { + return "" + } + switch seg { + case "admin", "healthcheck", "metrics": + return "" + } + return seg +} + +func firstNonEmptyHeader(headers http.Header, names ...string) string { + for _, name := range names { + if v := strings.TrimSpace(headers.Get(name)); v != "" { + return v + } + } + return "" +} + func normalizeSiweMessage(msg string) string { decoded, err := base64.StdEncoding.DecodeString(msg) if err != nil { diff --git a/auth/payload.go b/auth/payload.go index 712003534..9971ba7c5 100644 --- a/auth/payload.go +++ b/auth/payload.go @@ -3,11 +3,18 @@ package auth import "github.com/erpc/erpc/common" type AuthPayload struct { - Method string - Type common.AuthType - Secret *SecretPayload - Jwt *JwtPayload - Siwe *SiwePayload + Method string + Type common.AuthType + Secret *SecretPayload + Jwt *JwtPayload + Siwe *SiwePayload + ForwardedClientId *ForwardedClientIdPayload +} + +// ForwardedClientIdPayload carries a gateway-injected client id (not a secret). +type ForwardedClientIdPayload struct { + Value string + RateLimitBudget string } // This payload is used by both "secret" and "database" strategies diff --git a/auth/strategy_forwarded_client_id.go b/auth/strategy_forwarded_client_id.go new file mode 100644 index 000000000..af14caeea --- /dev/null +++ b/auth/strategy_forwarded_client_id.go @@ -0,0 +1,40 @@ +package auth + +import ( + "context" + "strings" + + "github.com/erpc/erpc/common" +) + +// ForwardedClientIdStrategy authenticates using a non-secret client id +// header injected by a trusted gateway after API-key verification +// (e.g. Envoy SecurityPolicy apiKeyAuth.forwardClientIDHeader). +type ForwardedClientIdStrategy struct { + cfg *common.ForwardedClientIdStrategyConfig +} + +var _ AuthStrategy = &ForwardedClientIdStrategy{} + +func NewForwardedClientIdStrategy(cfg *common.ForwardedClientIdStrategyConfig) *ForwardedClientIdStrategy { + return &ForwardedClientIdStrategy{cfg: cfg} +} + +func (s *ForwardedClientIdStrategy) Supports(ap *AuthPayload) bool { + return ap != nil && ap.Type == common.AuthTypeForwardedClientId +} + +func (s *ForwardedClientIdStrategy) Authenticate(ctx context.Context, req *common.NormalizedRequest, ap *AuthPayload) (*common.User, error) { + if ap == nil || ap.ForwardedClientId == nil || strings.TrimSpace(ap.ForwardedClientId.Value) == "" { + return nil, common.NewErrAuthUnauthorized("forwardedClientId", "missing client id header") + } + + id := strings.TrimSpace(ap.ForwardedClientId.Value) + user := &common.User{Id: id} + if s.cfg != nil && s.cfg.RateLimitBudget != "" { + user.RateLimitBudget = s.cfg.RateLimitBudget + } else if ap.ForwardedClientId.RateLimitBudget != "" { + user.RateLimitBudget = ap.ForwardedClientId.RateLimitBudget + } + return user, nil +} diff --git a/auth/strategy_forwarded_client_id_test.go b/auth/strategy_forwarded_client_id_test.go new file mode 100644 index 000000000..a6f987bbc --- /dev/null +++ b/auth/strategy_forwarded_client_id_test.go @@ -0,0 +1,104 @@ +package auth + +import ( + "context" + "net/http" + "net/url" + "testing" + + "github.com/erpc/erpc/common" + "github.com/stretchr/testify/require" +) + +func TestNewPayloadFromHttp_ForwardedClientId(t *testing.T) { + headers := http.Header{} + headers.Set("X-Client-Id", "cl-no-alpha") + ap, err := NewPayloadFromHttp("eth_blockNumber", "1.2.3.4:1234", headers, url.Values{}, "/") + require.NoError(t, err) + require.Equal(t, common.AuthTypeForwardedClientId, ap.Type) + require.NotNil(t, ap.ForwardedClientId) + require.Equal(t, "cl-no-alpha", ap.ForwardedClientId.Value) +} + +func TestNewPayloadFromHttp_PathSecret(t *testing.T) { + ap, err := NewPayloadFromHttp("eth_blockNumber", "1.2.3.4:1234", http.Header{}, url.Values{}, "/my-secret-key") + require.NoError(t, err) + require.Equal(t, common.AuthTypeSecret, ap.Type) + require.NotNil(t, ap.Secret) + require.Equal(t, "my-secret-key", ap.Secret.Value) + + // Trailing slash is cleaned to a single segment. + ap, err = NewPayloadFromHttp("eth_blockNumber", "1.2.3.4:1234", http.Header{}, url.Values{}, "/my-secret-key/") + require.NoError(t, err) + require.Equal(t, common.AuthTypeSecret, ap.Type) + require.Equal(t, "my-secret-key", ap.Secret.Value) +} + +func TestNewPayloadFromHttp_PathSecretIgnoredForMultiSegment(t *testing.T) { + ap, err := NewPayloadFromHttp("eth_blockNumber", "1.2.3.4:1234", http.Header{}, url.Values{}, "/main/evm/1") + require.NoError(t, err) + require.Equal(t, common.AuthTypeNetwork, ap.Type) +} + +func TestNewPayloadFromHttp_PathSecretIgnoredForReserved(t *testing.T) { + for _, seg := range []string{"/admin", "/healthcheck", "/metrics"} { + ap, err := NewPayloadFromHttp("eth_blockNumber", "1.2.3.4:1234", http.Header{}, url.Values{}, seg) + require.NoError(t, err) + require.Equal(t, common.AuthTypeNetwork, ap.Type, "path %s", seg) + } +} + +func TestSecretStrategy_RejectsEmpty(t *testing.T) { + s := NewSecretStrategy(&common.SecretStrategyConfig{Id: "cl-no-01", Value: ""}) + _, err := s.Authenticate(context.Background(), nil, &AuthPayload{ + Type: common.AuthTypeSecret, + Secret: &SecretPayload{Value: ""}, + }) + require.Error(t, err) + + s = NewSecretStrategy(&common.SecretStrategyConfig{Id: "cl-no-01", Value: "real-secret"}) + user, err := s.Authenticate(context.Background(), nil, &AuthPayload{ + Type: common.AuthTypeSecret, + Secret: &SecretPayload{Value: "real-secret"}, + }) + require.NoError(t, err) + require.Equal(t, "cl-no-01", user.Id) +} + +func TestNewPayloadFromHttp_ApiKeyQueryAndHeader(t *testing.T) { + ap, err := NewPayloadFromHttp("eth_blockNumber", "1.2.3.4:1234", http.Header{}, url.Values{"apikey": []string{"q-key"}}, "/") + require.NoError(t, err) + require.Equal(t, common.AuthTypeSecret, ap.Type) + require.Equal(t, "q-key", ap.Secret.Value) + + headers := http.Header{} + headers.Set("apikey", "h-key") + ap, err = NewPayloadFromHttp("eth_blockNumber", "1.2.3.4:1234", headers, url.Values{}, "/") + require.NoError(t, err) + require.Equal(t, common.AuthTypeSecret, ap.Type) + require.Equal(t, "h-key", ap.Secret.Value) +} + +func TestForwardedClientIdStrategy_Authenticate(t *testing.T) { + s := NewForwardedClientIdStrategy(&common.ForwardedClientIdStrategyConfig{ + Header: "X-Client-Id", + RateLimitBudget: "default-budget", + }) + ap := &AuthPayload{ + Type: common.AuthTypeForwardedClientId, + ForwardedClientId: &ForwardedClientIdPayload{ + Value: "cl-no-beta", + }, + } + user, err := s.Authenticate(context.Background(), nil, ap) + require.NoError(t, err) + require.Equal(t, "cl-no-beta", user.Id) + require.Equal(t, "default-budget", user.RateLimitBudget) +} + +func TestForwardedClientIdStrategy_MissingHeader(t *testing.T) { + s := NewForwardedClientIdStrategy(&common.ForwardedClientIdStrategyConfig{}) + ap := &AuthPayload{Type: common.AuthTypeForwardedClientId} + _, err := s.Authenticate(context.Background(), nil, ap) + require.Error(t, err) +} diff --git a/auth/strategy_secret.go b/auth/strategy_secret.go index 271e18a26..51dc9026e 100644 --- a/auth/strategy_secret.go +++ b/auth/strategy_secret.go @@ -2,6 +2,7 @@ package auth import ( "context" + "strings" "github.com/erpc/erpc/common" ) @@ -21,6 +22,14 @@ func (s *SecretStrategy) Supports(ap *AuthPayload) bool { } func (s *SecretStrategy) Authenticate(ctx context.Context, req *common.NormalizedRequest, ap *AuthPayload) (*common.User, error) { + if ap == nil || ap.Secret == nil { + return nil, common.NewErrAuthUnauthorized("secret", "missing secret") + } + // Reject empty configured or presented secrets so a missing env expansion + // (value="") can never authenticate an empty path/query credential. + if strings.TrimSpace(s.cfg.Value) == "" || strings.TrimSpace(ap.Secret.Value) == "" { + return nil, common.NewErrAuthUnauthorized("secret", "invalid secret") + } if ap.Secret.Value != s.cfg.Value { return nil, common.NewErrAuthUnauthorized("secret", "invalid secret") } diff --git a/clients/registry.go b/clients/registry.go index 533d140ec..2175019c1 100644 --- a/clients/registry.go +++ b/clients/registry.go @@ -95,7 +95,9 @@ func (manager *ClientRegistry) CreateClient(appCtx context.Context, ups common.U creation.once.Do(func() { lg := manager.logger.With().Str("upstreamId", cfg.Id).Logger() switch cfg.Type { - case common.UpstreamTypeEvm: + case common.UpstreamTypeEvm, common.UpstreamTypeJsonRpc: + // jsonrpc architecture reuses the generic HTTP/WS JSON-RPC clients + // (passthrough; no EVM chainId/state poller). gRPC BDS remains EVM-only. if parsedUrl.Scheme == "http" || parsedUrl.Scheme == "https" { newClient, err = NewGenericHttpJsonRpcClient( appCtx, @@ -123,6 +125,8 @@ func (manager *ClientRegistry) CreateClient(appCtx context.Context, ups common.U if err != nil { clientErr = fmt.Errorf("failed to create WebSocket client for upstream %v: %w", cfg.Id, err) } + } else if (parsedUrl.Scheme == "grpc" || parsedUrl.Scheme == "grpc+bds") && cfg.Type == common.UpstreamTypeJsonRpc { + clientErr = fmt.Errorf("unsupported endpoint scheme: %v for upstream type jsonrpc: %v", parsedUrl.Scheme, cfg.Id) } else if parsedUrl.Scheme == "grpc" || parsedUrl.Scheme == "grpc+bds" { grpcPoolSize := 0 if cfg.Grpc != nil { diff --git a/common/architecture_jsonrpc.go b/common/architecture_jsonrpc.go new file mode 100644 index 000000000..dcd37b471 --- /dev/null +++ b/common/architecture_jsonrpc.go @@ -0,0 +1,56 @@ +package common + +import ( + "context" + "net/http" +) + +const ( + UpstreamTypeJsonRpc UpstreamType = "jsonrpc" +) + +// JsonRpcNetworkConfig identifies a non-EVM JSON-RPC network (Solana, Starknet, …). +// Network id becomes jsonrpc:. No chainId / state poller / EVM method hooks. +type JsonRpcNetworkConfig struct { + // Id is a stable slug (usually the CLL alias), e.g. solana-mainnet. + Id string `yaml:"id" json:"id"` +} + +func init() { + RegisterArchitecture(ArchitectureJsonRpc, &JsonRpcArchitectureHandler{}) +} + +// JsonRpcArchitectureHandler is the passthrough handler for jsonrpc networks: +// no architecture hooks. Its error extractor claims nothing, so upstream errors +// fall through to the EVM extractor in the composite (as before the registry). +type JsonRpcArchitectureHandler struct{} + +func (h *JsonRpcArchitectureHandler) HandleProjectPreForward(ctx context.Context, network Network, req *NormalizedRequest) (bool, *NormalizedResponse, error) { + return false, nil, nil +} + +func (h *JsonRpcArchitectureHandler) HandleNetworkPreForward(ctx context.Context, network Network, upstreams []Upstream, req *NormalizedRequest) (bool, *NormalizedResponse, error) { + return false, nil, nil +} + +func (h *JsonRpcArchitectureHandler) HandleNetworkPostForward(ctx context.Context, network Network, req *NormalizedRequest, resp *NormalizedResponse, err error) (*NormalizedResponse, error) { + return resp, err +} + +func (h *JsonRpcArchitectureHandler) HandleUpstreamPreForward(ctx context.Context, network Network, upstream Upstream, req *NormalizedRequest, skipCacheRead bool) (bool, *NormalizedResponse, error) { + return false, nil, nil +} + +func (h *JsonRpcArchitectureHandler) HandleUpstreamPostForward(ctx context.Context, network Network, upstream Upstream, req *NormalizedRequest, resp *NormalizedResponse, err error, skipCacheRead bool) (*NormalizedResponse, error) { + return resp, err +} + +func (h *JsonRpcArchitectureHandler) NewJsonRpcErrorExtractor() JsonRpcErrorExtractor { + return noopJsonRpcErrorExtractor{} +} + +type noopJsonRpcErrorExtractor struct{} + +func (noopJsonRpcErrorExtractor) Extract(*http.Response, *NormalizedResponse, *JsonRpcResponse, Upstream) error { + return nil +} diff --git a/common/architecture_jsonrpc_test.go b/common/architecture_jsonrpc_test.go new file mode 100644 index 000000000..3047e3b25 --- /dev/null +++ b/common/architecture_jsonrpc_test.go @@ -0,0 +1,63 @@ +package common + +import ( + "testing" + + "gopkg.in/yaml.v3" +) + +func TestJsonRpcNetworkId(t *testing.T) { + n := &NetworkConfig{ + Architecture: ArchitectureJsonRpc, + Alias: "solana-mainnet", + JsonRpc: &JsonRpcNetworkConfig{Id: "solana-mainnet"}, + } + if got := n.NetworkId(); got != "jsonrpc:solana-mainnet" { + t.Fatalf("NetworkId()=%q", got) + } + if !IsValidArchitecture(string(ArchitectureJsonRpc)) { + t.Fatal("ArchitectureJsonRpc should be valid") + } + if !IsValidNetwork("jsonrpc:solana-mainnet") { + t.Fatal("jsonrpc:solana-mainnet should be valid") + } + if IsValidNetwork("jsonrpc:") || IsValidNetwork("jsonrpc:a:b") { + t.Fatal("invalid jsonrpc ids accepted") + } +} + +func TestJsonRpcNetworkConfig_UnmarshalYAML_FailsafeObject(t *testing.T) { + // Chart emits failsafe as a single object (not a list). NetworkConfig must + // still accept architecture/jsonRpc via the oldNetworkConfig fallback. + const raw = ` +architecture: jsonrpc +alias: solana-mainnet +jsonRpc: + id: solana-mainnet +failsafe: + timeout: + duration: 30s + retry: + maxAttempts: 5 + delay: 50ms +` + var n NetworkConfig + if err := yaml.Unmarshal([]byte(raw), &n); err != nil { + t.Fatalf("unmarshal: %v", err) + } + if n.Architecture != ArchitectureJsonRpc { + t.Fatalf("architecture=%q", n.Architecture) + } + if n.JsonRpc == nil || n.JsonRpc.Id != "solana-mainnet" { + t.Fatalf("jsonRpc=%v", n.JsonRpc) + } + if n.Alias != "solana-mainnet" { + t.Fatalf("alias=%q", n.Alias) + } + if len(n.Failsafe) != 1 || n.Failsafe[0].Timeout == nil { + t.Fatalf("failsafe not converted from object: %+v", n.Failsafe) + } + if got := n.NetworkId(); got != "jsonrpc:solana-mainnet" { + t.Fatalf("NetworkId()=%q", got) + } +} diff --git a/common/config.go b/common/config.go index 0ddc48864..9c0b5f416 100644 --- a/common/config.go +++ b/common/config.go @@ -1309,6 +1309,9 @@ type JsonRpcUpstreamConfig struct { EnableGzip *bool `yaml:"enableGzip,omitempty" json:"enableGzip"` Headers map[string]string `yaml:"headers,omitempty" json:"headers"` ProxyPool string `yaml:"proxyPool,omitempty" json:"proxyPool"` + // NetworkId binds this upstream to architecture jsonrpc (slug only, e.g. solana-mainnet). + // Required when type is jsonrpc; skipped for EVM upstreams that already use jsonRpc for batch/headers. + NetworkId string `yaml:"networkId,omitempty" json:"networkId,omitempty"` } func (c *JsonRpcUpstreamConfig) Copy() *JsonRpcUpstreamConfig { @@ -2301,6 +2304,7 @@ type NetworkConfig struct { Failsafe []*FailsafeConfig `yaml:"failsafe,omitempty" json:"failsafe"` Evm *EvmNetworkConfig `yaml:"evm,omitempty" json:"evm"` Svm *SvmNetworkConfig `yaml:"svm,omitempty" json:"svm"` + JsonRpc *JsonRpcNetworkConfig `yaml:"jsonRpc,omitempty" json:"jsonRpc"` SelectionPolicy *SelectionPolicyConfig `yaml:"selectionPolicy,omitempty" json:"selectionPolicy"` DirectiveDefaults *DirectiveDefaultsConfig `yaml:"directiveDefaults,omitempty" json:"directiveDefaults"` Alias string `yaml:"alias,omitempty" json:"alias"` @@ -2376,6 +2380,7 @@ func (n *NetworkConfig) UnmarshalYAML(unmarshal func(interface{}) error) error { Failsafe *FailsafeConfig `yaml:"failsafe,omitempty"` Evm *EvmNetworkConfig `yaml:"evm,omitempty"` Svm *SvmNetworkConfig `yaml:"svm,omitempty"` + JsonRpc *JsonRpcNetworkConfig `yaml:"jsonRpc,omitempty"` SelectionPolicy *SelectionPolicyConfig `yaml:"selectionPolicy,omitempty"` DirectiveDefaults *DirectiveDefaultsConfig `yaml:"directiveDefaults,omitempty"` Alias string `yaml:"alias,omitempty"` @@ -2395,6 +2400,7 @@ func (n *NetworkConfig) UnmarshalYAML(unmarshal func(interface{}) error) error { n.RateLimitBudget = old.RateLimitBudget n.Evm = old.Evm n.Svm = old.Svm + n.JsonRpc = old.JsonRpc n.SelectionPolicy = old.SelectionPolicy n.DirectiveDefaults = old.DirectiveDefaults n.Alias = old.Alias @@ -2885,11 +2891,12 @@ func (s *SelectionPolicyConfig) UnmarshalYAML(unmarshal func(interface{}) error) type AuthType string const ( - AuthTypeSecret AuthType = "secret" - AuthTypeDatabase AuthType = "database" - AuthTypeJwt AuthType = "jwt" - AuthTypeSiwe AuthType = "siwe" - AuthTypeNetwork AuthType = "network" + AuthTypeSecret AuthType = "secret" + AuthTypeDatabase AuthType = "database" + AuthTypeJwt AuthType = "jwt" + AuthTypeSiwe AuthType = "siwe" + AuthTypeNetwork AuthType = "network" + AuthTypeForwardedClientId AuthType = "forwardedClientId" ) type AuthConfig struct { @@ -2920,12 +2927,27 @@ type AuthStrategyConfig struct { // NormalizedRequest.SetUserFromTrustedHeader). AllowClientDirectives *string `yaml:"allowClientDirectives,omitempty" json:"allowClientDirectives,omitempty"` - Type AuthType `yaml:"type" json:"type" tstype:"TsAuthType"` - Network *NetworkStrategyConfig `yaml:"network,omitempty" json:"network,omitempty"` - Secret *SecretStrategyConfig `yaml:"secret,omitempty" json:"secret,omitempty"` - Database *DatabaseStrategyConfig `yaml:"database,omitempty" json:"database,omitempty"` - Jwt *JwtStrategyConfig `yaml:"jwt,omitempty" json:"jwt,omitempty"` - Siwe *SiweStrategyConfig `yaml:"siwe,omitempty" json:"siwe,omitempty"` + Type AuthType `yaml:"type" json:"type" tstype:"TsAuthType"` + Network *NetworkStrategyConfig `yaml:"network,omitempty" json:"network,omitempty"` + Secret *SecretStrategyConfig `yaml:"secret,omitempty" json:"secret,omitempty"` + Database *DatabaseStrategyConfig `yaml:"database,omitempty" json:"database,omitempty"` + Jwt *JwtStrategyConfig `yaml:"jwt,omitempty" json:"jwt,omitempty"` + Siwe *SiweStrategyConfig `yaml:"siwe,omitempty" json:"siwe,omitempty"` + ForwardedClientId *ForwardedClientIdStrategyConfig `yaml:"forwardedClientId,omitempty" json:"forwardedClientId,omitempty"` +} + +// ForwardedClientIdStrategyConfig trusts a non-secret client identity header +// injected by an upstream gateway after API-key auth (e.g. Envoy +// apiKeyAuth.forwardClientIDHeader → X-Client-Id). Must only be enabled +// behind a gateway that overwrites/strips client-supplied values of that header. +type ForwardedClientIdStrategyConfig struct { + // Header documents the expected gateway identity header (default conceptually + // "X-Client-Id"). Payload extraction in auth.NewPayloadFromHttp currently + // always reads X-Client-Id via case-insensitive Header.Get; this field is + // not yet used to select the header name at runtime. + Header string `yaml:"header,omitempty" json:"header,omitempty"` + // RateLimitBudget, if set, is applied to the authenticated user. + RateLimitBudget string `yaml:"rateLimitBudget,omitempty" json:"rateLimitBudget,omitempty"` } type SecretStrategyConfig struct { @@ -3222,6 +3244,11 @@ func (c *NetworkConfig) NetworkId() string { return "" } return util.SvmNetworkId(c.Svm.Chain, c.Svm.Cluster) + case ArchitectureJsonRpc: + if c.JsonRpc == nil || c.JsonRpc.Id == "" { + return "" + } + return util.JsonRpcNetworkId(c.JsonRpc.Id) default: return "" } diff --git a/common/defaults.go b/common/defaults.go index c312f5ba5..037149f4c 100644 --- a/common/defaults.go +++ b/common/defaults.go @@ -1938,8 +1938,17 @@ func (u *UpstreamConfig) SetDefaults(defaults *UpstreamConfig) error { } } if u.Type == "" { - // TODO make actual calls to detect other types (solana, btc, etc)? - u.Type = UpstreamTypeEvm + if u.JsonRpc != nil && u.JsonRpc.NetworkId != "" { + u.Type = UpstreamTypeJsonRpc + } else { + // TODO make actual calls to detect other types (solana, btc, etc)? + u.Type = UpstreamTypeEvm + } + } + if u.Type == UpstreamTypeJsonRpc { + if u.JsonRpc == nil { + u.JsonRpc = &JsonRpcUpstreamConfig{} + } } if len(u.Failsafe) > 0 { @@ -2297,18 +2306,26 @@ func (n *NetworkConfig) SetDefaults(upstreams []*UpstreamConfig, defaults *Netwo if n.Architecture == "" { if n.Evm != nil { - n.Architecture = "evm" + n.Architecture = ArchitectureEvm } else if n.Svm != nil { n.Architecture = ArchitectureSvm + } else if n.JsonRpc != nil && n.JsonRpc.Id != "" { + n.Architecture = ArchitectureJsonRpc } } - if n.Architecture == "evm" && n.Evm == nil { + if n.Architecture == ArchitectureEvm && n.Evm == nil { n.Evm = &EvmNetworkConfig{} } if n.Architecture == ArchitectureSvm && n.Svm == nil { n.Svm = &SvmNetworkConfig{} } + if n.Architecture == ArchitectureJsonRpc && n.JsonRpc == nil { + n.JsonRpc = &JsonRpcNetworkConfig{} + } + if n.Architecture == ArchitectureJsonRpc && n.JsonRpc.Id == "" && n.Alias != "" { + n.JsonRpc.Id = n.Alias + } // Apply methods defaults if n.Methods == nil { @@ -3249,6 +3266,23 @@ func (s *AuthStrategyConfig) SetDefaults() error { } } + if s.Type == AuthTypeForwardedClientId && s.ForwardedClientId == nil { + s.ForwardedClientId = &ForwardedClientIdStrategyConfig{} + } + if s.ForwardedClientId != nil { + s.Type = AuthTypeForwardedClientId + if err := s.ForwardedClientId.SetDefaults(); err != nil { + return fmt.Errorf("failed to set defaults for forwardedClientId strategy: %w", err) + } + } + + return nil +} + +func (s *ForwardedClientIdStrategyConfig) SetDefaults() error { + if s.Header == "" { + s.Header = "X-Client-Id" + } return nil } diff --git a/common/json_rpc.go b/common/json_rpc.go index 636d46346..b38cd286e 100644 --- a/common/json_rpc.go +++ b/common/json_rpc.go @@ -373,7 +373,10 @@ func (r *JsonRpcResponse) ParseFromStream(ctx []context.Context, reader io.Reade } } - if len(temp.Error) > 0 { + // Bitcoin-family nodes often emit `"error": null` on success. Treat that + // (and missing error) as no error — ParseError("null") used to invent a + // server-side exception and fail the whole upstream attempt. + if len(temp.Error) > 0 && string(temp.Error) != "null" { if err := r.ParseError(string(temp.Error)); err != nil { return err } @@ -401,11 +404,17 @@ func (r *JsonRpcResponse) ParseError(raw string) error { r.errBytes = nil + // JSON-RPC allows "error": null (Bitcoin Core / dogecoind / litecoind). + // That means success — do not fabricate a server-side exception. + if raw == "null" { + return nil + } + // First attempt to unmarshal the error as a typical JSON-RPC error var rpcErr ErrJsonRpcExceptionExternal if err := SonicCfg.UnmarshalFromString(raw, &rpcErr); err != nil { // Special case: check for non-standard error structures in the raw data - if raw == "" || raw == "null" { + if raw == "" { r.Error = NewErrJsonRpcExceptionExternal( int(JsonRpcErrorServerSideException), "unexpected empty response from upstream endpoint", @@ -1204,10 +1213,13 @@ func isEmptyishValue(v interface{}) bool { type JsonRpcRequest struct { sync.RWMutex - JSONRPC string `json:"jsonrpc,omitempty"` - ID interface{} `json:"id,omitempty"` - Method string `json:"method"` - Params []interface{} `json:"params"` + JSONRPC string `json:"jsonrpc,omitempty"` + ID interface{} `json:"id,omitempty"` + Method string `json:"method"` + // omitempty: Stellar (and some other non-EVM) reject "params":[] — they + // expect the field absent (or an object). Empty/nil params are omitted on + // the wire; EVM nodes accept both forms. + Params []interface{} `json:"params,omitempty"` // idRaw stores the verbatim bytes of the id as received from the client. // This is used to round-trip the id back without precision loss for ids diff --git a/common/json_rpc_test.go b/common/json_rpc_test.go index 1a78528c4..542236699 100644 --- a/common/json_rpc_test.go +++ b/common/json_rpc_test.go @@ -220,10 +220,23 @@ func TestJsonRpcRequest_MarshalParams(t *testing.T) { }) assert.NoError(t, err) - expectedRawReq := `{"jsonrpc":"2.0","id":1,"method":"eth_blockNumber","params":[]}` + // Empty params omitted — required for Stellar / some non-EVM nodes + // that reject "params":[] (expect absent field or object). + expectedRawReq := `{"jsonrpc":"2.0","id":1,"method":"eth_blockNumber"}` assert.Equal(t, expectedRawReq, string(rawReq)) }) + t.Run("Nil", func(t *testing.T) { + rawReq, err := SonicCfg.Marshal(JsonRpcRequest{ + JSONRPC: "2.0", + ID: 1, + Method: "getLatestLedger", + Params: nil, + }) + assert.NoError(t, err) + assert.Equal(t, `{"jsonrpc":"2.0","id":1,"method":"getLatestLedger"}`, string(rawReq)) + }) + t.Run("Value", func(t *testing.T) { rawReq, err := SonicCfg.Marshal(JsonRpcRequest{ JSONRPC: "2.0", @@ -238,6 +251,20 @@ func TestJsonRpcRequest_MarshalParams(t *testing.T) { }) } +func TestJsonRpcResponse_ErrorNullIsSuccess(t *testing.T) { + // Bitcoin-family (dogecoind/litecoind) returns "error":null on success. + raw := `{"result":67851510,"error":null,"id":1}` + r := &JsonRpcResponse{} + err := r.ParseFromStream(nil, bytes.NewReader([]byte(raw)), len(raw)) + require.NoError(t, err) + assert.Nil(t, r.Error) + assert.Equal(t, "67851510", string(r.result)) + + r2 := &JsonRpcResponse{} + require.NoError(t, r2.ParseError("null")) + assert.Nil(t, r2.Error) +} + func TestJsonRpcResponse_CanonicalHash_EmptyishNormalization(t *testing.T) { // Test cases that should produce the same hash due to emptyish normalization testGroups := []struct { diff --git a/common/network.go b/common/network.go index 8b1367f11..742257c4a 100644 --- a/common/network.go +++ b/common/network.go @@ -12,8 +12,9 @@ import ( type NetworkArchitecture string const ( - ArchitectureEvm NetworkArchitecture = "evm" - ArchitectureSvm NetworkArchitecture = "svm" + ArchitectureEvm NetworkArchitecture = "evm" + ArchitectureSvm NetworkArchitecture = "svm" + ArchitectureJsonRpc NetworkArchitecture = "jsonrpc" ) type Network interface { @@ -87,7 +88,7 @@ func EvmLeaderUpstream(n Network, ctx context.Context) Upstream { func IsValidArchitecture(architecture string) bool { switch NetworkArchitecture(architecture) { - case ArchitectureEvm, ArchitectureSvm: + case ArchitectureEvm, ArchitectureSvm, ArchitectureJsonRpc: return true } return false @@ -143,6 +144,10 @@ func IsValidNetwork(network string) bool { } return isIdentifier(chain) && isIdentifier(cluster) } + if strings.HasPrefix(network, "jsonrpc:") { + id := strings.TrimPrefix(network, "jsonrpc:") + return id != "" && !strings.Contains(id, ":") + } return false } diff --git a/common/request.go b/common/request.go index a68358e14..544624c20 100644 --- a/common/request.go +++ b/common/request.go @@ -271,6 +271,10 @@ type NormalizedRequest struct { // Resolved client IP (set by HTTP ingress using trusted forwarders) clientIP atomic.Value + // Client transport that delivered this request ("http" or "ws"). + // Defaults to "http" when unset so HTTP ingress needs no explicit set. + transport atomic.Value + // Per-request execution counters; lazy-init via execStateHolder. execStateHolder execStateHolder } @@ -1310,6 +1314,27 @@ func (r *NormalizedRequest) AgentName() string { return "unknown" } +// SetTransport records the client ingress transport ("http" or "ws"). +func (r *NormalizedRequest) SetTransport(transport string) { + if r == nil || transport == "" { + return + } + r.transport.Store(transport) +} + +// Transport returns the client ingress transport. Defaults to "http". +func (r *NormalizedRequest) Transport() string { + if r == nil { + return "http" + } + if v := r.transport.Load(); v != nil { + if s, ok := v.(string); ok && s != "" { + return s + } + } + return "http" +} + // getUserAgent returns the user agent string, with query parameter taking precedence over header func (r *NormalizedRequest) getUserAgent(headers http.Header, queryArgs url.Values) string { // Query parameter takes precedence diff --git a/common/validation.go b/common/validation.go index 9cc9ae1ef..bd30c805a 100644 --- a/common/validation.go +++ b/common/validation.go @@ -912,6 +912,13 @@ func (s *AuthStrategyConfig) Validate() error { if err := s.Database.Validate(); err != nil { return err } + case AuthTypeForwardedClientId: + if s.ForwardedClientId == nil { + return fmt.Errorf("auth.*.forwardedClientId is required for forwardedClientId strategy") + } + if err := s.ForwardedClientId.Validate(); err != nil { + return err + } default: return fmt.Errorf("auth.*.type '%s' is invalid must be one of: %v", s.Type, []AuthType{ AuthTypeNetwork, @@ -919,11 +926,19 @@ func (s *AuthStrategyConfig) Validate() error { AuthTypeJwt, AuthTypeSiwe, AuthTypeDatabase, + AuthTypeForwardedClientId, }) } return nil } +func (s *ForwardedClientIdStrategyConfig) Validate() error { + if s == nil { + return fmt.Errorf("auth.*.forwardedClientId is required") + } + return nil +} + func (s *DatabaseStrategyConfig) Validate() error { if s.Connector == nil { return fmt.Errorf("auth.*.database.connector is required") @@ -952,7 +967,13 @@ func (s *NetworkStrategyConfig) Validate() error { } func (s *SecretStrategyConfig) Validate() error { - if s.Value == "" { + if s == nil { + return fmt.Errorf("auth.*.secret is required") + } + if strings.TrimSpace(s.Id) == "" { + return fmt.Errorf("auth.*.secret.id is required") + } + if strings.TrimSpace(s.Value) == "" { return fmt.Errorf("auth.*.secret.value is required") } return nil @@ -1036,6 +1057,14 @@ func (u *UpstreamConfig) Validate(c *Config, skipEndpointCheck bool) error { if !skipEndpointCheck && u.Endpoint == "" { return fmt.Errorf("upstream.*.endpoint is required") } + if u.Type == UpstreamTypeJsonRpc { + if u.JsonRpc == nil || u.JsonRpc.NetworkId == "" { + return fmt.Errorf("upstream.*.jsonRpc.networkId is required for type jsonrpc") + } + if !util.IsValidIdentifier(u.JsonRpc.NetworkId) { + return fmt.Errorf("upstream.*.jsonRpc.networkId '%s' is invalid", u.JsonRpc.NetworkId) + } + } if u.Evm != nil { if err := u.Evm.Validate(u); err != nil { return err @@ -1496,12 +1525,20 @@ func (n *NetworkConfig) Validate(c *Config) error { if n.Architecture == "" { return fmt.Errorf("network.*.architecture is required") } - if n.Architecture == "evm" && n.Evm == nil { + if n.Architecture == ArchitectureEvm && n.Evm == nil { return fmt.Errorf("network.*.evm is required for evm networks") } if n.Architecture == ArchitectureSvm && n.Svm == nil { return fmt.Errorf("network.*.svm is required for svm networks") } + if n.Architecture == ArchitectureJsonRpc { + if n.JsonRpc == nil || n.JsonRpc.Id == "" { + return fmt.Errorf("network.*.jsonRpc.id is required for jsonrpc networks") + } + if !util.IsValidIdentifier(n.JsonRpc.Id) { + return fmt.Errorf("network.*.jsonRpc.id '%s' must contain only alphanumeric characters, dash, or underscore", n.JsonRpc.Id) + } + } if n.Evm != nil { if err := n.Evm.Validate(); err != nil { return err diff --git a/erpc/healthcheck.go b/erpc/healthcheck.go index 5349b547a..283d4f9bf 100644 --- a/erpc/healthcheck.go +++ b/erpc/healthcheck.go @@ -86,7 +86,7 @@ func (s *HttpServer) handleHealthCheck( headers := r.Header queryArgs := r.URL.Query() - ap, err := auth.NewPayloadFromHttp("healthcheck", r.RemoteAddr, headers, queryArgs) + ap, err := auth.NewPayloadFromHttp("healthcheck", r.RemoteAddr, headers, queryArgs, r.URL.Path) if err != nil { handleErrorResponse(ctx, &logger, startedAt, nil, err, w, encoder, writeFatalError, &common.TRUE, s.executionHeadersMode()) return diff --git a/erpc/http_server.go b/erpc/http_server.go index 94aacb4c0..a59ebdee9 100644 --- a/erpc/http_server.go +++ b/erpc/http_server.go @@ -581,9 +581,9 @@ func (s *HttpServer) createRequestHandler() http.Handler { var ap *auth.AuthPayload if project != nil { - ap, err = auth.NewPayloadFromHttp(method, r.RemoteAddr, headers, queryArgs) + ap, err = auth.NewPayloadFromHttp(method, r.RemoteAddr, headers, queryArgs, r.URL.Path) } else if isAdmin { - ap, err = auth.NewPayloadFromHttp(method, r.RemoteAddr, headers, queryArgs) + ap, err = auth.NewPayloadFromHttp(method, r.RemoteAddr, headers, queryArgs, r.URL.Path) } if err != nil { responses[index] = processErrorBody(&rlg, &startedAt, nq, err, &common.TRUE) diff --git a/erpc/http_server_test.go b/erpc/http_server_test.go index ba47c0a00..1163317a9 100644 --- a/erpc/http_server_test.go +++ b/erpc/http_server_test.go @@ -4136,7 +4136,7 @@ func TestHttpServer_HandleHealthCheck(t *testing.T) { pp.networksRegistry = NewNetworksRegistry(pp, ctx, pp.upstreamsRegistry, mtk, nil, nil, nil, nil, logger) authReg, _ := auth.NewAuthRegistry(ctx, logger, "test", &common.AuthConfig{Strategies: []*common.AuthStrategyConfig{ - {Type: common.AuthTypeSecret, Secret: &common.SecretStrategyConfig{Value: "test-secret"}}, + {Type: common.AuthTypeSecret, Secret: &common.SecretStrategyConfig{Id: "test-user", Value: "test-secret"}}, }}, nil) return &HttpServer{ diff --git a/erpc/networks.go b/erpc/networks.go index d1ad6ebd6..6aaae4a63 100644 --- a/erpc/networks.go +++ b/erpc/networks.go @@ -2770,9 +2770,9 @@ func (n *Network) prepareRequest(ctx context.Context, nr *common.NormalizedReque ) } evm.NormalizeHttpJsonRpc(ctx, nr, jsonRpcReq) - case common.ArchitectureSvm: - // SVM doesn't need any EVM-style normalization (hex padding, block tag expansion, etc.). - // Validate that the request parses as JSON-RPC and move on. + case common.ArchitectureSvm, common.ArchitectureJsonRpc: + // SVM doesn't need any EVM-style normalization (hex padding, block tag expansion, etc.) + // and jsonrpc is a passthrough. Validate that the request parses as JSON-RPC and move on. if _, err := nr.JsonRpcRequest(ctx); err != nil { return common.NewErrJsonRpcExceptionInternal( 0, diff --git a/erpc/networks_registry.go b/erpc/networks_registry.go index 5b3f75018..45a37fd1e 100644 --- a/erpc/networks_registry.go +++ b/erpc/networks_registry.go @@ -339,6 +339,8 @@ func (nr *NetworksRegistry) prepareNetwork(nwCfg *common.NetworkConfig) (*Networ if nr.svmJsonRpcCache != nil { network.cacheDal = nr.svmJsonRpcCache.WithProjectId(nr.project.Config.Id) } + case common.ArchitectureJsonRpc: + // No architecture-specific cache yet; failsafe + metrics still apply. } // Register alias for lazy-created networks to support alias-based routing if nwCfg.Alias != "" { @@ -401,6 +403,11 @@ func (nr *NetworksRegistry) resolveNetworkConfig(networkId string) (*common.Netw return nil, common.NewErrInvalidEvmChainId(networkId) } nwCfg.Svm = &common.SvmNetworkConfig{Chain: chain, Cluster: cluster} + case common.ArchitectureJsonRpc: + if !util.IsValidIdentifier(s[1]) { + return nil, fmt.Errorf("invalid jsonrpc network id: %s", networkId) + } + nwCfg.JsonRpc = &common.JsonRpcNetworkConfig{Id: s[1]} default: return nil, common.NewErrInvalidEvmChainId(networkId) } diff --git a/erpc/ws_server.go b/erpc/ws_server.go index af2fd9cbe..5bd66fa83 100644 --- a/erpc/ws_server.go +++ b/erpc/ws_server.go @@ -306,7 +306,7 @@ func (wsc *WsConnection) handleRequest(ctx context.Context, raw []byte, startedA return unsupported("subscription methods (eth_subscribe, eth_unsubscribe) are not supported in batch requests") } - ap, err := auth.NewPayloadFromHttp(method, wsc.httpReq.RemoteAddr, headers, queryArgs) + ap, err := auth.NewPayloadFromHttp(method, wsc.httpReq.RemoteAddr, headers, queryArgs, wsc.httpReq.URL.Path) if err != nil { return fail(err, &common.TRUE) } diff --git a/typescript/config/src/generated.ts b/typescript/config/src/generated.ts index 88b36abad..232fda8e2 100644 --- a/typescript/config/src/generated.ts +++ b/typescript/config/src/generated.ts @@ -1884,6 +1884,7 @@ export const AuthTypeDatabase: AuthType = "database"; export const AuthTypeJwt: AuthType = "jwt"; export const AuthTypeSiwe: AuthType = "siwe"; export const AuthTypeNetwork: AuthType = "network"; +export const AuthTypeForwardedClientId: AuthType = "forwardedClientId"; export interface AuthConfig { strategies: TsAuthStrategyConfig[]; } @@ -1915,6 +1916,17 @@ export interface AuthStrategyConfig { database?: DatabaseStrategyConfig; jwt?: JwtStrategyConfig; siwe?: SiweStrategyConfig; + /** + * Trust a gateway-injected client id header (e.g. Envoy X-Client-Id). + */ + forwardedClientId?: ForwardedClientIdStrategyConfig; +} +export interface ForwardedClientIdStrategyConfig { + /** + * Header carrying the client id. Default: "X-Client-Id". + */ + header?: string; + rateLimitBudget?: string; } export interface SecretStrategyConfig { id: string; diff --git a/typescript/config/src/index.ts b/typescript/config/src/index.ts index 0c7da1a48..7497cd36a 100644 --- a/typescript/config/src/index.ts +++ b/typescript/config/src/index.ts @@ -71,6 +71,7 @@ export { AuthTypeJwt, AuthTypeSiwe, AuthTypeNetwork, + AuthTypeForwardedClientId, // Consensus related ConsensusLowParticipantsBehaviorReturnError, ConsensusLowParticipantsBehaviorAcceptMostCommonValidResult, diff --git a/upstream/registry.go b/upstream/registry.go index a77a9273b..0452deafa 100644 --- a/upstream/registry.go +++ b/upstream/registry.go @@ -520,6 +520,8 @@ func (u *UpstreamsRegistry) buildUpstreamBootstrapTask(upsCfg *common.UpstreamCo taskName := fmt.Sprintf("upstream/%s", cfg.Id) if cfg.Evm != nil && cfg.Evm.ChainId > 0 { taskName = fmt.Sprintf("network/%s/upstream/%s", util.EvmNetworkId(cfg.Evm.ChainId), cfg.Id) + } else if cfg.Type == common.UpstreamTypeJsonRpc && cfg.JsonRpc != nil && cfg.JsonRpc.NetworkId != "" { + taskName = fmt.Sprintf("network/%s/upstream/%s", util.JsonRpcNetworkId(cfg.JsonRpc.NetworkId), cfg.Id) } return util.NewBootstrapTask( taskName, diff --git a/upstream/upstream.go b/upstream/upstream.go index 251dce30d..509134ed9 100644 --- a/upstream/upstream.go +++ b/upstream/upstream.go @@ -1413,6 +1413,14 @@ func (u *Upstream) detectFeatures(ctx context.Context) error { // Genesis-hash validation runs in Bootstrap (svmVerifyGenesisHash) once // the client and networkId are in place, so it can go through the // upstream's normal Forward path. + } else if cfg.Type == common.UpstreamTypeJsonRpc { + if cfg.JsonRpc == nil || cfg.JsonRpc.NetworkId == "" { + return common.NewTaskFatal(fmt.Errorf("upstream.*.jsonRpc.networkId is required for type jsonrpc")) + } + if !util.IsValidIdentifier(cfg.JsonRpc.NetworkId) { + return common.NewTaskFatal(fmt.Errorf("upstream.*.jsonRpc.networkId '%s' is invalid", cfg.JsonRpc.NetworkId)) + } + u.networkId.Store(util.JsonRpcNetworkId(cfg.JsonRpc.NetworkId)) } else { return fmt.Errorf("upstream type not supported: %s", cfg.Type) } diff --git a/util/ids.go b/util/ids.go index f9016bf61..bd3a0f94e 100644 --- a/util/ids.go +++ b/util/ids.go @@ -24,6 +24,10 @@ func SvmNetworkId(chain, cluster string) string { return "svm:" + chain + ":" + cluster } +func JsonRpcNetworkId(id string) string { + return fmt.Sprintf("jsonrpc:%s", id) +} + var validIdentifierRegex = regexp.MustCompile(`^[a-zA-Z0-9_-]+$`) func IsValidIdentifier(s string) bool { @@ -63,6 +67,10 @@ func IsValidNetworkId(s string) bool { } return true } + if strings.HasPrefix(s, "jsonrpc:") { + id := s[len("jsonrpc:"):] + return id != "" && IsValidIdentifier(id) + } return false } From 11108fd204dada413e2545eb84acc67c809cf8fc Mon Sep 17 00:00:00 2001 From: snowkide Date: Fri, 2 Oct 2026 11:50:58 +0200 Subject: [PATCH 2/2] feat(metrics): count every WS call and expose per-client WS and latency metrics - network_request_received_total gains transport (http|ws); WS ingress sets it on single and batch requests and on eth_(un)subscribe - ws_requests_total{category,user,agent_name,outcome}: every JSON-RPC request over WebSocket counted once, including requests rejected before the network (invalid, method_not_allowed, unauthorized, rate_limited, network_unavailable, invalid_batch). category stays "n/a" until the client is authenticated so unauthenticated clients cannot mint series - ws_subscription_events_total{user,agent_name}: notifications delivered - websocket_subscription_notifications_dropped_total gains user and agent_name (labels frozen at subscribe time, no network lookup per drop) - ws_subscriptions_active{kind,user}, ws_connections_active, ws_connections_closed_total{close_code,initiator} - client_request_duration_seconds{network,user,transport,outcome}: per-user end-to-end latency without the category/vendor/upstream dimensions, for deployments that drop user from network_request_duration_seconds Co-Authored-By: Claude Opus 5.5 --- erpc/projects.go | 26 ++++++- erpc/subscription_manager.go | 17 ++++- erpc/ws_server.go | 65 ++++++++++++++++-- erpc/ws_server_metrics_test.go | 83 +++++++++++++++++++++++ indexer/adapters/wsclient/adapter.go | 48 ++++++++++++- indexer/adapters/wsclient/adapter_test.go | 48 ++++++++++--- telemetry/metrics.go | 56 ++++++++++++++- 7 files changed, 322 insertions(+), 21 deletions(-) create mode 100644 erpc/ws_server_metrics_test.go diff --git a/erpc/projects.go b/erpc/projects.go index cbcd5a905..cd976a2a0 100644 --- a/erpc/projects.go +++ b/erpc/projects.go @@ -128,6 +128,30 @@ func (p *PreparedProject) AuthenticateConsumer(ctx context.Context, req *common. } func (p *PreparedProject) Forward(ctx context.Context, networkId string, nq *common.NormalizedRequest) (*common.NormalizedResponse, error) { + start := time.Now() + resp, err := p.forward(ctx, networkId, nq) + if nw := nq.Network(); nw != nil { + telemetry.ObserverHandle(telemetry.MetricClientRequestDuration, + p.Config.Id, nw.Label(), nq.UserId(), nq.Transport(), clientRequestOutcome(err), + ).Observe(time.Since(start).Seconds()) + } + return resp, err +} + +// clientRequestOutcome is the outcome label of client_request_duration_seconds +// and ws_requests_total for a request that reached the project. +func clientRequestOutcome(err error) string { + switch { + case err == nil: + return "ok" + case common.HasErrorCode(err, common.ErrCodeProjectRateLimitRuleExceeded, common.ErrCodeAuthRateLimitRuleExceeded): + return "rate_limited" + default: + return "error" + } +} + +func (p *PreparedProject) forward(ctx context.Context, networkId string, nq *common.NormalizedRequest) (*common.NormalizedResponse, error) { start := time.Now() ctx, span := common.StartDetailSpan(ctx, "Project.Forward") defer span.End() @@ -150,7 +174,7 @@ func (p *PreparedProject) Forward(ctx context.Context, networkId string, nq *com reqFinality := nq.Finality(ctx) telemetry.CounterHandle(telemetry.MetricNetworkRequestsReceived, - p.Config.Id, network.Label(), method, reqFinality.String(), nq.UserId(), nq.AgentName(), + p.Config.Id, network.Label(), method, reqFinality.String(), nq.UserId(), nq.AgentName(), nq.Transport(), ).Inc() lg := p.Logger.With(). Str("component", "proxy"). diff --git a/erpc/subscription_manager.go b/erpc/subscription_manager.go index 49b5c1dff..eafe28c1b 100644 --- a/erpc/subscription_manager.go +++ b/erpc/subscription_manager.go @@ -105,7 +105,7 @@ func (sm *SubscriptionManager) Subscribe( reqFinality := nq.Finality(ctx) telemetry.CounterHandle(telemetry.MetricNetworkRequestsReceived, project.Config.Id, nw.Label(), method, - reqFinality.String(), nq.UserId(), nq.AgentName(), + reqFinality.String(), nq.UserId(), nq.AgentName(), nq.Transport(), ).Inc() jrReq, err := nq.JsonRpcRequest() @@ -128,7 +128,12 @@ func (sm *SubscriptionManager) Subscribe( return nil, err } - if err := conn.adapter.AddSubscription(clientSubID, networkId, kind, filterHash, maxSubs); err != nil { + if err := conn.adapter.AddSubscription(clientSubID, networkId, kind, filterHash, maxSubs, wsclient.SubscriptionLabels{ + Project: project.Config.Id, + Network: nw.Label(), + User: nq.UserId(), + AgentName: nq.AgentName(), + }); err != nil { sm.releaseFilter(ctx, networkId, kind, filterHash) if errors.Is(err, wsclient.ErrLimitExceeded) { err = common.NewErrSubscriptionLimitExceeded(maxSubs) @@ -171,7 +176,7 @@ func (sm *SubscriptionManager) Unsubscribe( reqFinality := nq.Finality(ctx) telemetry.CounterHandle(telemetry.MetricNetworkRequestsReceived, project.Config.Id, nw.Label(), method, - reqFinality.String(), nq.UserId(), nq.AgentName(), + reqFinality.String(), nq.UserId(), nq.AgentName(), nq.Transport(), ).Inc() jrReq, err := nq.JsonRpcRequest() @@ -406,6 +411,9 @@ func (sm *SubscriptionManager) recordSuccessMetrics( project.Config.Id, nw.Label(), "proxy", "proxy", method, finality.String(), nq.UserId(), ).Observe(time.Since(start).Seconds()) + telemetry.ObserverHandle(telemetry.MetricClientRequestDuration, + project.Config.Id, nw.Label(), nq.UserId(), nq.Transport(), "ok", + ).Observe(time.Since(start).Seconds()) } func (sm *SubscriptionManager) recordFailureMetrics( @@ -430,6 +438,9 @@ func (sm *SubscriptionManager) recordFailureMetrics( project.Config.Id, nw.Label(), "", "", method, finality.String(), nq.UserId(), ).Observe(time.Since(start).Seconds()) + telemetry.ObserverHandle(telemetry.MetricClientRequestDuration, + project.Config.Id, nw.Label(), nq.UserId(), nq.Transport(), "error", + ).Observe(time.Since(start).Seconds()) } // networkHandle adapts *Network to indexer.NetworkHandle. diff --git a/erpc/ws_server.go b/erpc/ws_server.go index 5bd66fa83..3577bcb1b 100644 --- a/erpc/ws_server.go +++ b/erpc/ws_server.go @@ -9,6 +9,7 @@ import ( "io" "net/http" "runtime/debug" + "strconv" "sync" "sync/atomic" "time" @@ -54,6 +55,13 @@ type WsConnection struct { closed atomic.Bool // no more writes once set peerClosed atomic.Bool // the peer's close frame was already answered + // peerEndCode is the close code the read loop ended on when the peer + // ended the connection (1006 when it vanished without a close frame); + // 0 when the server initiated the close. + peerEndCode atomic.Int32 + + // networkLabel is the metrics label of networkId (alias if configured). + networkLabel string // subsClosed is guarded by SubscriptionManager.connMu. subsClosed bool @@ -105,6 +113,12 @@ func (s *HttpServer) handleWebSocket( done: make(chan struct{}), } + wsc.networkLabel = wsc.networkId + if nw, err := project.GetNetwork(r.Context(), wsc.networkId); err == nil { + wsc.networkLabel = nw.Label() + } + telemetry.GaugeHandle(telemetry.MetricWsConnectionsActive, project.Config.Id, wsc.networkLabel).Inc() + lg.Info().Str("connId", wsc.id).Str("remoteAddr", r.RemoteAddr).Msg("websocket connection established") // Store before checking draining so shutdownWebSockets either sees this @@ -209,6 +223,9 @@ func (wsc *WsConnection) readLoop(pongWait time.Duration) { Str("errType", fmt.Sprintf("%T", err)) if ce, ok := err.(*websocket.CloseError); ok { ev = ev.Int("closeCode", ce.Code).Str("closeReason", ce.Text) + wsc.peerEndCode.Store(int32(ce.Code)) + } else { + wsc.peerEndCode.Store(websocket.CloseAbnormalClosure) } if websocket.IsUnexpectedCloseError(err, websocket.CloseNormalClosure, websocket.CloseGoingAway) { ev.Msg("websocket read ended: unexpected close") @@ -270,7 +287,16 @@ func (wsc *WsConnection) handleMessage(raw []byte) { // Subscription methods are rejected in batches. func (wsc *WsConnection) handleRequest(ctx context.Context, raw []byte, startedAt *time.Time, inBatch bool) interface{} { nq := common.NewNormalizedRequest(raw) + nq.SetTransport("ws") nq.ForwardHeaders = make(http.Header) + + // Every request is counted exactly once, including those rejected before + // they reach the network. category stays "n/a" until the client is + // authenticated so it cannot mint series with arbitrary method names. + outcome, category := "invalid", "n/a" + defer func() { + wsc.countRequest(category, nq.UserId(), nq.AgentName(), outcome) + }() requestCtx := common.StartRequestSpan(ctx, nq) nq.SetClientIP(wsc.server.resolveRealClientIP(wsc.httpReq)) @@ -300,25 +326,33 @@ func (wsc *WsConnection) handleRequest(ctx context.Context, raw []byte, startedA return fail(err, &common.TRUE) } if !allowed { + outcome = "method_not_allowed" return unsupported(fmt.Sprintf("method not supported: %s", method)) } if inBatch && IsSubscriptionMethod(method) { + outcome = "method_not_allowed" return unsupported("subscription methods (eth_subscribe, eth_unsubscribe) are not supported in batch requests") } + outcome = "unauthorized" ap, err := auth.NewPayloadFromHttp(method, wsc.httpReq.RemoteAddr, headers, queryArgs, wsc.httpReq.URL.Path) if err != nil { return fail(err, &common.TRUE) } user, err := project.AuthenticateConsumer(requestCtx, nq, method, ap) if err != nil { + if common.HasErrorCode(err, common.ErrCodeAuthRateLimitRuleExceeded) { + outcome = "rate_limited" + } return fail(err, wsc.server.serverCfg.IncludeErrorDetails) } nq.SetUser(user) if project.Config.TrustUserIdHeader && nq.User() == nil { nq.SetUserFromTrustedHeader(headers.Get(common.HeaderUserId)) } + category = method + outcome = "network_unavailable" nw, err := project.GetNetwork(requestCtx, wsc.networkId) if err != nil { return fail(err, wsc.server.serverCfg.IncludeErrorDetails) @@ -335,6 +369,7 @@ func (wsc *WsConnection) handleRequest(ctx context.Context, raw []byte, startedA default: resp, err = project.Forward(requestCtx, wsc.networkId, nq) } + outcome = clientRequestOutcome(err) if err != nil { if resp != nil { go resp.Release() @@ -349,11 +384,13 @@ func (wsc *WsConnection) handleRequest(ctx context.Context, raw []byte, startedA func (wsc *WsConnection) handleBatch(ctx context.Context, raw []byte, startedAt *time.Time) { var requests []json.RawMessage if err := common.SonicCfg.Unmarshal(raw, &requests); err != nil { + wsc.countRequest("n/a", "", "", "invalid_batch") wsc.writeError(int(common.JsonRpcErrorParseException), "parse error") return } // JSON-RPC 2.0 section 6: an empty batch gets a single Invalid Request error. if len(requests) == 0 { + wsc.countRequest("n/a", "", "", "invalid_batch") wsc.writeError(int(common.JsonRpcErrorClientSideException), "invalid request: empty batch") return } @@ -532,6 +569,15 @@ func (wsc *WsConnection) teardown() { wsc.server.activeWsConns.Delete(wsc.id) close(wsc.done) + initiator, code := "server", wsc.closeCode + if pc := wsc.peerEndCode.Load(); pc != 0 { + initiator, code = "client", int(pc) + } + telemetry.CounterHandle(telemetry.MetricWsConnectionsClosedTotal, + wsc.project.Config.Id, wsc.networkLabel, strconv.Itoa(code), initiator, + ).Inc() + telemetry.GaugeHandle(telemetry.MetricWsConnectionsActive, wsc.project.Config.Id, wsc.networkLabel).Dec() + wsc.logger.Info().Str("connId", wsc.id).Int("closeCode", wsc.closeCode).Msg("websocket connection closed") } @@ -564,12 +610,8 @@ func (wsc *WsConnection) WriteSubscriptionNotification(clientSubId string, resul // the client reconnects knowing it missed data instead of silently // diverging. func (wsc *WsConnection) NotificationDropped(sub wsclient.Subscription, lossy bool) { - network := sub.NetworkID - if nw, err := wsc.project.GetNetwork(wsc.ctx, sub.NetworkID); err == nil { - network = nw.Label() - } telemetry.CounterHandle(telemetry.MetricWebsocketSubscriptionNotificationsDroppedTotal, - wsc.project.Config.Id, network, sub.Kind.String(), + wsc.project.Config.Id, sub.Labels.Network, sub.Kind.String(), sub.Labels.User, sub.Labels.AgentName, ).Inc() now := time.Now().UnixNano() @@ -583,3 +625,16 @@ func (wsc *WsConnection) NotificationDropped(sub wsclient.Subscription, lossy bo wsc.stop(websocket.CloseTryAgainLater, "subscription buffer overflow: notifications were dropped", 0) } } + +// countRequest increments ws_requests_total for one client request. +func (wsc *WsConnection) countRequest(category, user, agent, outcome string) { + if user == "" { + user = "n/a" + } + if agent == "" { + agent = "unknown" + } + telemetry.CounterHandle(telemetry.MetricWsRequestsTotal, + wsc.project.Config.Id, wsc.networkLabel, category, user, agent, outcome, + ).Inc() +} diff --git a/erpc/ws_server_metrics_test.go b/erpc/ws_server_metrics_test.go new file mode 100644 index 000000000..f2eef6f95 --- /dev/null +++ b/erpc/ws_server_metrics_test.go @@ -0,0 +1,83 @@ +package erpc + +import ( + "testing" + "time" + + "github.com/erpc/erpc/telemetry" + "github.com/erpc/erpc/util" + "github.com/gorilla/websocket" + "github.com/prometheus/client_golang/prometheus/testutil" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// wsRequests reads ws_requests_total. The agent is only known once the +// request is authenticated (the Go dialer's User-Agent maps to "go"). +func wsRequests(category, agent, outcome string) float64 { + return testutil.ToFloat64(telemetry.CounterHandle(telemetry.MetricWsRequestsTotal, + "test_ws", "evm:123", category, "n/a", agent, outcome)) +} + +// TestWebSocket_RequestAccounting: every client request over WebSocket is +// counted once in ws_requests_total, including those rejected before they +// reach the network, and forwarded ones carry transport="ws" on +// network_request_received_total. +func TestWebSocket_RequestAccounting(t *testing.T) { + setupGock() + defer util.ResetGock() + + cfg := httpOnlyConfig() + cfg.Projects[0].IgnoreMethods = []string{"debug_*"} + addr, cleanup := setupTestERPCServer(t, cfg) + defer cleanup() + + activeConns := func() float64 { + return testutil.ToFloat64(telemetry.GaugeHandle(telemetry.MetricWsConnectionsActive, "test_ws", "evm:123")) + } + closedByClient := func() float64 { + return testutil.ToFloat64(telemetry.CounterHandle(telemetry.MetricWsConnectionsClosedTotal, + "test_ws", "evm:123", "1000", "client")) + } + received := func() float64 { + return testutil.ToFloat64(telemetry.CounterHandle(telemetry.MetricNetworkRequestsReceived, + "test_ws", "evm:123", "eth_getBalance", "realtime", "n/a", "go", "ws")) + } + + okBefore := wsRequests("eth_getBalance", "go", "ok") + notAllowedBefore := wsRequests("n/a", "unknown", "method_not_allowed") + invalidBefore := wsRequests("n/a", "unknown", "invalid") + invalidBatchBefore := wsRequests("n/a", "unknown", "invalid_batch") + receivedBefore := received() + activeBefore := activeConns() + closedBefore := closedByClient() + + conn := dialWs(t, addr) + require.Eventually(t, func() bool { return activeConns() == activeBefore+1 }, 2*time.Second, 10*time.Millisecond) + + resp := sendAndReceive(t, conn, `{"jsonrpc":"2.0","id":1,"method":"eth_getBalance","params":["0x1234567890abcdef1234567890abcdef12345678","latest"]}`) + assert.Equal(t, "0xabc123", resp["result"]) + + resp = sendAndReceive(t, conn, `{"jsonrpc":"2.0","id":2,"method":"debug_traceTransaction","params":["0xabc"]}`) + assert.NotNil(t, resp["error"]) + + resp = sendAndReceive(t, conn, `{"jsonrpc":"2.0","id":3}`) + assert.NotNil(t, resp["error"]) + + resp = sendAndReceive(t, conn, `[{"jsonrpc":"2.0"`) + assert.NotNil(t, resp["error"]) + + assert.Equal(t, okBefore+1, wsRequests("eth_getBalance", "go", "ok")) + assert.Equal(t, notAllowedBefore+1, wsRequests("n/a", "unknown", "method_not_allowed")) + assert.Equal(t, invalidBefore+1, wsRequests("n/a", "unknown", "invalid")) + assert.Equal(t, invalidBatchBefore+1, wsRequests("n/a", "unknown", "invalid_batch")) + assert.Equal(t, receivedBefore+1, received(), "forwarded WS request must carry transport=ws") + + require.NoError(t, conn.WriteControl(websocket.CloseMessage, + websocket.FormatCloseMessage(websocket.CloseNormalClosure, ""), time.Now().Add(time.Second))) + conn.Close() + + require.Eventually(t, func() bool { + return activeConns() == activeBefore && closedByClient() == closedBefore+1 + }, 5*time.Second, 20*time.Millisecond) +} diff --git a/indexer/adapters/wsclient/adapter.go b/indexer/adapters/wsclient/adapter.go index 73889c78f..34e8d7a6e 100644 --- a/indexer/adapters/wsclient/adapter.go +++ b/indexer/adapters/wsclient/adapter.go @@ -8,6 +8,8 @@ import ( "sync" "github.com/erpc/erpc/indexer" + "github.com/erpc/erpc/telemetry" + "github.com/prometheus/client_golang/prometheus" "github.com/rs/zerolog" ) @@ -56,11 +58,35 @@ type routeKey struct { filterHash string } +// SubscriptionLabels are frozen at eth_subscribe time and label the +// per-client metrics of a subscription (delivered, dropped, active). +type SubscriptionLabels struct { + Project string + // Network is the metrics network label (alias if configured, else id). + Network string + User string + AgentName string +} + +func (l SubscriptionLabels) withDefaults() SubscriptionLabels { + if l.Project == "" { + l.Project = "n/a" + } + if l.User == "" { + l.User = "n/a" + } + if l.AgentName == "" { + l.AgentName = "unknown" + } + return l +} + type clientSub struct { id string kind indexer.EventKind networkID string filterHash string // "" for newHeads + labels SubscriptionLabels notify chan json.RawMessage done chan struct{} // closed by whoever removes the sub from Adapter.subs @@ -111,12 +137,17 @@ func (a *Adapter) Deliver(ev indexer.IndexedEvent) { // AddSubscription registers a subscription and starts its writer. It fails // with ErrLimitExceeded when the connection already holds maxSubs, and with // ErrClosed once the adapter was drained. -func (a *Adapter) AddSubscription(clientSubID, networkID string, kind indexer.EventKind, filterHash string, maxSubs int) error { +func (a *Adapter) AddSubscription(clientSubID, networkID string, kind indexer.EventKind, filterHash string, maxSubs int, labels SubscriptionLabels) error { + labels = labels.withDefaults() + if labels.Network == "" { + labels.Network = networkID + } sub := &clientSub{ id: clientSubID, kind: kind, networkID: networkID, filterHash: filterHash, + labels: labels, notify: make(chan json.RawMessage, a.bufferSize), done: make(chan struct{}), } @@ -140,6 +171,7 @@ func (a *Adapter) AddSubscription(clientSubID, networkID string, kind indexer.Ev set[clientSubID] = struct{}{} a.mu.Unlock() + sub.activeGauge().Inc() go a.runWriter(sub) return nil } @@ -164,6 +196,7 @@ func (a *Adapter) RemoveSubscription(clientSubID string) (kind indexer.EventKind a.mu.Unlock() close(sub.done) + sub.activeGauge().Dec() return sub.kind, sub.networkID, sub.filterHash, true } @@ -181,6 +214,7 @@ func (a *Adapter) Drain() []Subscription { out := make([]Subscription, 0, len(subs)) for _, sub := range subs { close(sub.done) + sub.activeGauge().Dec() out = append(out, sub.snapshot()) } return out @@ -192,15 +226,22 @@ func (sub *clientSub) snapshot() Subscription { Kind: sub.kind, NetworkID: sub.networkID, FilterHash: sub.filterHash, + Labels: sub.labels, } } +func (sub *clientSub) activeGauge() prometheus.Gauge { + return telemetry.GaugeHandle(telemetry.MetricWsSubscriptionsActive, + sub.labels.Project, sub.labels.Network, sub.kind.String(), sub.labels.User) +} + // Subscription describes one subscription. type Subscription struct { ClientSubID string Kind indexer.EventKind NetworkID string FilterHash string + Labels SubscriptionLabels } // Count returns the number of active subscriptions. @@ -219,7 +260,12 @@ func (a *Adapter) runWriter(sub *clientSub) { if err := a.writer.WriteSubscriptionNotification(sub.id, payload); err != nil { a.logger.Debug().Err(err).Str("clientSubId", sub.id). Msg("failed to write subscription notification") + continue } + telemetry.CounterHandle(telemetry.MetricWsSubscriptionEventsTotal, + sub.labels.Project, sub.labels.Network, sub.kind.String(), + sub.labels.User, sub.labels.AgentName, + ).Inc() } } } diff --git a/indexer/adapters/wsclient/adapter_test.go b/indexer/adapters/wsclient/adapter_test.go index 09960b559..d6210620e 100644 --- a/indexer/adapters/wsclient/adapter_test.go +++ b/indexer/adapters/wsclient/adapter_test.go @@ -9,6 +9,8 @@ import ( "time" "github.com/erpc/erpc/indexer" + "github.com/erpc/erpc/telemetry" + "github.com/prometheus/client_golang/prometheus/testutil" "github.com/rs/zerolog" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -58,13 +60,13 @@ func deliver(a *Adapter, kind indexer.EventKind, filterHash string, n int) { func TestAdapter_AddAfterDrainIsRefused(t *testing.T) { a := newTestAdapter(&fakeWriter{}, 8) - require.NoError(t, a.AddSubscription("s1", "evm:1", indexer.KindLog, "f", 10)) + require.NoError(t, a.AddSubscription("s1", "evm:1", indexer.KindLog, "f", 10, SubscriptionLabels{})) removed := a.Drain() require.Len(t, removed, 1) assert.Equal(t, "s1", removed[0].ClientSubID) - assert.ErrorIs(t, a.AddSubscription("s2", "evm:1", indexer.KindLog, "f", 10), ErrClosed) + assert.ErrorIs(t, a.AddSubscription("s2", "evm:1", indexer.KindLog, "f", 10, SubscriptionLabels{}), ErrClosed) assert.Equal(t, 0, a.Count()) } @@ -72,8 +74,8 @@ func TestAdapter_AddAfterDrainIsRefused(t *testing.T) { // its filter reference is released exactly once. func TestAdapter_RemoveAndDrainAreExclusive(t *testing.T) { a := newTestAdapter(&fakeWriter{}, 8) - require.NoError(t, a.AddSubscription("s1", "evm:1", indexer.KindLog, "f", 10)) - require.NoError(t, a.AddSubscription("s2", "evm:1", indexer.KindLog, "f", 10)) + require.NoError(t, a.AddSubscription("s1", "evm:1", indexer.KindLog, "f", 10, SubscriptionLabels{})) + require.NoError(t, a.AddSubscription("s2", "evm:1", indexer.KindLog, "f", 10, SubscriptionLabels{})) _, _, _, existed := a.RemoveSubscription("s1") require.True(t, existed) @@ -93,7 +95,7 @@ func TestAdapter_SubscriptionLimitIsAtomic(t *testing.T) { wg.Add(1) go func(i int) { defer wg.Done() - if a.AddSubscription(fmt.Sprintf("s%d", i), "evm:1", indexer.KindNewHead, "", 10) == nil { + if a.AddSubscription(fmt.Sprintf("s%d", i), "evm:1", indexer.KindNewHead, "", 10, SubscriptionLabels{}) == nil { ok.Add(1) } }(i) @@ -108,7 +110,7 @@ func TestAdapter_SubscriptionLimitIsAtomic(t *testing.T) { func TestAdapter_BurstWithinBufferIsDelivered(t *testing.T) { w := &fakeWriter{} a := newTestAdapter(w, 256) - require.NoError(t, a.AddSubscription("s1", "evm:1", indexer.KindLog, "f", 10)) + require.NoError(t, a.AddSubscription("s1", "evm:1", indexer.KindLog, "f", 10, SubscriptionLabels{})) defer a.Drain() deliver(a, indexer.KindLog, "f", 200) @@ -125,7 +127,7 @@ func TestAdapter_BurstWithinBufferIsDelivered(t *testing.T) { func TestAdapter_NewHeadsOverflowDropsOldest(t *testing.T) { w := &fakeWriter{block: make(chan struct{})} a := newTestAdapter(w, 2) - require.NoError(t, a.AddSubscription("s1", "evm:1", indexer.KindNewHead, "", 10)) + require.NoError(t, a.AddSubscription("s1", "evm:1", indexer.KindNewHead, "", 10, SubscriptionLabels{})) defer a.Drain() // The writer takes one event and blocks on it; two more fill the buffer. @@ -148,7 +150,7 @@ func TestAdapter_NewHeadsOverflowDropsOldest(t *testing.T) { func TestAdapter_LogsOverflowIsReportedNotEvicted(t *testing.T) { w := &fakeWriter{block: make(chan struct{})} a := newTestAdapter(w, 2) - require.NoError(t, a.AddSubscription("s1", "evm:1", indexer.KindLog, "f", 10)) + require.NoError(t, a.AddSubscription("s1", "evm:1", indexer.KindLog, "f", 10, SubscriptionLabels{})) defer a.Drain() deliver(a, indexer.KindLog, "f", 1) @@ -164,3 +166,33 @@ func TestAdapter_LogsOverflowIsReportedNotEvicted(t *testing.T) { assert.Equal(t, []string{"0", "0", "1"}, written, "queued logs must be kept in order") assert.Equal(t, []bool{false, false, false}, dropped) } + +// Delivered notifications and the active-subscription gauge are labeled +// with the labels frozen at subscribe time. +func TestAdapter_PerClientMetrics(t *testing.T) { + labels := SubscriptionLabels{Project: "p-metrics", Network: "eth-mainnet", User: "cl-no-99", AgentName: "go"} + delivered := func() float64 { + return testutil.ToFloat64(telemetry.CounterHandle(telemetry.MetricWsSubscriptionEventsTotal, + "p-metrics", "eth-mainnet", indexer.KindNewHead.String(), "cl-no-99", "go")) + } + active := func() float64 { + return testutil.ToFloat64(telemetry.GaugeHandle(telemetry.MetricWsSubscriptionsActive, + "p-metrics", "eth-mainnet", indexer.KindNewHead.String(), "cl-no-99")) + } + + deliveredBefore := delivered() + w := &fakeWriter{} + a := newTestAdapter(w, 8) + require.NoError(t, a.AddSubscription("s1", "evm:1", indexer.KindNewHead, "", 10, labels)) + require.NoError(t, a.AddSubscription("s2", "evm:1", indexer.KindNewHead, "", 10, labels)) + assert.Equal(t, float64(2), active()) + + deliver(a, indexer.KindNewHead, "", 3) + require.Eventually(t, func() bool { return delivered() == deliveredBefore+6 }, 2*time.Second, 10*time.Millisecond) + + _, _, _, existed := a.RemoveSubscription("s1") + require.True(t, existed) + assert.Equal(t, float64(1), active()) + a.Drain() + assert.Equal(t, float64(0), active()) +} diff --git a/telemetry/metrics.go b/telemetry/metrics.go index 033c525e5..cdb029302 100644 --- a/telemetry/metrics.go +++ b/telemetry/metrics.go @@ -123,11 +123,51 @@ var ( Help: "Whether the upstream WebSocket connection is currently established (1) or down/wedged (0).", }, []string{"project", "vendor", "network", "upstream"}) - MetricWebsocketSubscriptionNotificationsDroppedTotal = DefineCounter(prometheus.CounterOpts{ + MetricWebsocketSubscriptionNotificationsDroppedTotal = DefineLabeledCounter(prometheus.CounterOpts{ Namespace: "erpc", Name: "websocket_subscription_notifications_dropped_total", Help: "Subscription notifications dropped because a client's per-subscription buffer was full. For logs the client connection is closed with 1013.", - }, []string{"project", "network", "kind"}) + }, []string{"project", "network", "kind", "user", "agent_name"}) + + // MetricWsSubscriptionEventsTotal counts subscription notifications + // (newHeads / logs / pending txs) written to a client. These are server + // pushes, not client calls: client JSON-RPC calls over WebSocket are in + // network_request_received_total{transport="ws"} and ws_requests_total. + MetricWsSubscriptionEventsTotal = DefineLabeledCounter(prometheus.CounterOpts{ + Namespace: "erpc", + Name: "ws_subscription_events_total", + Help: "Subscription notifications successfully written to a client WebSocket.", + }, []string{"project", "network", "kind", "user", "agent_name"}) + + // MetricWsRequestsTotal counts every JSON-RPC request a client sends over + // WebSocket (batch items individually), including those rejected before + // they reach the network (invalid, method not allowed, unauthorized, rate + // limited) and therefore absent from network_request_received_total. + // category is "n/a" until the request is authenticated, so unauthenticated + // clients cannot mint series with arbitrary method names. + MetricWsRequestsTotal = DefineLabeledCounter(prometheus.CounterOpts{ + Namespace: "erpc", + Name: "ws_requests_total", + Help: "JSON-RPC requests received over client WebSocket connections, by outcome (ok, error, invalid, method_not_allowed, unauthorized, rate_limited, network_unavailable, invalid_batch).", + }, []string{"project", "network", "category", "user", "agent_name", "outcome"}) + + MetricWsConnectionsActive = DefineGauge(prometheus.GaugeOpts{ + Namespace: "erpc", + Name: "ws_connections_active", + Help: "Client WebSocket connections currently open.", + }, []string{"project", "network"}) + + MetricWsConnectionsClosedTotal = DefineCounter(prometheus.CounterOpts{ + Namespace: "erpc", + Name: "ws_connections_closed_total", + Help: "Client WebSocket connections closed, by close code and whether the peer initiated the close.", + }, []string{"project", "network", "close_code", "initiator"}) + + MetricWsSubscriptionsActive = DefineGauge(prometheus.GaugeOpts{ + Namespace: "erpc", + Name: "ws_subscriptions_active", + Help: "Client subscriptions currently active.", + }, []string{"project", "network", "kind", "user"}) // MetricNetworkServedTipBlockNumber is the block number the network actually // advertises/serves as the tip for a block tag (axis=latest|finalized): the @@ -645,7 +685,7 @@ var ( Namespace: "erpc", Name: "network_request_received_total", Help: "Total number of requests received for a network.", - }, []string{"project", "network", "category", "finality", "user", "agent_name"}) + }, []string{"project", "network", "category", "finality", "user", "agent_name", "transport"}) MetricNetworkMultiplexedRequests = DefineLabeledCounter(prometheus.CounterOpts{ Namespace: "erpc", @@ -1070,6 +1110,16 @@ var ( Help: "Duration of requests for a network.", }, []string{"project", "network", "vendor", "upstream", "category", "finality", "user"}) + // MetricClientRequestDuration is end-to-end request duration per client + // (user) without the category/vendor/upstream dimensions, so per-user + // latency stays affordable where network_request_duration_seconds has its + // user label dropped. + MetricClientRequestDuration = DefineLabeledHistogram(prometheus.HistogramOpts{ + Namespace: "erpc", + Name: "client_request_duration_seconds", + Help: "End-to-end duration of client requests by network, user and transport.", + }, []string{"project", "network", "user", "transport", "outcome"}) + // Per-request integrity latency overhead — the time a request waited on // integrity data-checks plus aux force-fetches (canonical header/receipts), // summed across attempts; excludes failover latency from rejections. Uses the