diff --git a/pkg/chipingress/client.go b/pkg/chipingress/client.go index 6065e52f64..f51ce93bb7 100644 --- a/pkg/chipingress/client.go +++ b/pkg/chipingress/client.go @@ -5,6 +5,7 @@ import ( "crypto/tls" "fmt" "net" + "strings" "time" "github.com/google/uuid" @@ -112,14 +113,22 @@ func NewClient(address string, opts ...Opt) (Client, error) { if cfg.perRPCCredentials != nil { grpcOpts = append(grpcOpts, grpc.WithPerRPCCredentials(cfg.perRPCCredentials)) } - // Add headers as a unary interceptor, use for non-auth headers + // Add headers as unary interceptors, use for non-auth headers. + // WithChainUnaryInterceptor is used (rather than WithUnaryInterceptor) so that + // headerProvider and nopInfoHeaderProvider compose instead of the second call + // silently overriding the first (grpc.WithUnaryInterceptor is last-one-wins). + var unaryInterceptors []grpc.UnaryClientInterceptor if cfg.headerProvider != nil { - grpcOpts = append(grpcOpts, grpc.WithUnaryInterceptor(newHeaderInterceptor(cfg.headerProvider))) + unaryInterceptors = append(unaryInterceptors, newHeaderInterceptor(cfg.headerProvider)) // NOTE: not supporting streaming interceptors } if cfg.nopInfoHeaderProvider != nil { - grpcOpts = append(grpcOpts, grpc.WithUnaryInterceptor(newHeaderInterceptor(cfg.nopInfoHeaderProvider))) + unaryInterceptors = append(unaryInterceptors, newHeaderInterceptor(cfg.nopInfoHeaderProvider)) + } + + if len(unaryInterceptors) > 0 { + grpcOpts = append(grpcOpts, grpc.WithChainUnaryInterceptor(unaryInterceptors...)) } conn, err := grpc.NewClient(address, grpcOpts...) @@ -206,6 +215,13 @@ func WithHeaderProvider(provider HeaderProvider) Opt { return func(c *clientConfig) { c.headerProvider = provider } } +// WithResourceAttributeHeaders returns an Opt that attaches the provided resource attributes +// as sanitized gRPC metadata headers. It combines SanitizeMetadataHeaders with +// NewStaticHeaderProvider so the safe, validated path is used by default. +func WithResourceAttributeHeaders(attrs map[string]string) Opt { + return WithHeaderProvider(NewStaticHeaderProvider(SanitizeMetadataHeaders(attrs))) +} + // WithInsecureConnection configures the client to use an insecure connection (no TLS). func WithInsecureConnection() Opt { return func(config *clientConfig) { @@ -267,9 +283,42 @@ func newHeaderInterceptor(provider HeaderProvider) grpc.UnaryClientInterceptor { } } +// EventOpt configures a CloudEvent after its well-known attributes have been set by NewEvent. +type EventOpt func(*ce.Event) + +// sanitizeExtensionName lower-cases name and strips every rune outside [a-z0-9], the character +// set the CloudEvents spec requires for extension attribute names. +func sanitizeExtensionName(name string) string { + var b strings.Builder + for _, r := range strings.ToLower(name) { + if (r >= 'a' && r <= 'z') || (r >= '0' && r <= '9') { + b.WriteRune(r) + } + } + return b.String() +} + +// WithResourceAttributeExtensions returns an EventOpt that sets a CloudEvent extension for each +// entry in attrs, sanitizing keys via sanitizeExtensionName so they satisfy the CloudEvents +// extension-name character set. Entries that sanitize to an empty string, or that collide with a +// reserved extension name (see reservedExtensionNames), are skipped. Keys are applied in sorted +// order so that if two distinct keys sanitize to the same name, the result is deterministic. +func WithResourceAttributeExtensions(attrs map[string]string) EventOpt { + return func(event *ce.Event) { + for _, pair := range sanitizeResourceAttributeKeys(attrs, nil) { + event.SetExtension(pair.name, attrs[pair.key]) + } + } +} + // NewEvent creates a new CloudEvent with the specified domain, entity, payload, and optional attributes. func NewEvent(domain, entity string, payload []byte, attributes map[string]any) (CloudEvent, error) { + return NewEventWithOpts(domain, entity, payload, attributes) +} +// NewEventWithOpts creates a new CloudEvent like NewEvent, additionally applying opts (e.g. +// WithResourceAttributeExtensions) to the event before its data is set. +func NewEventWithOpts(domain, entity string, payload []byte, attributes map[string]any, opts ...EventOpt) (CloudEvent, error) { event := ce.NewEvent() event.SetSource(domain) event.SetType(entity) @@ -303,6 +352,10 @@ func NewEvent(domain, entity string, payload []byte, attributes map[string]any) event.SetExtension(IdempotencyKeyAttr, val) } + for _, opt := range opts { + opt(&event) + } + err := event.SetData(ceformat.ContentTypeProtobuf, payload) if err != nil { return ce.Event{}, fmt.Errorf("could not set data on event: %w", err) diff --git a/pkg/chipingress/client_test.go b/pkg/chipingress/client_test.go index c1a21c588b..c224e1bb42 100644 --- a/pkg/chipingress/client_test.go +++ b/pkg/chipingress/client_test.go @@ -3,6 +3,7 @@ package chipingress import ( "context" "fmt" + "net" "testing" "time" @@ -167,6 +168,89 @@ func TestNewEvent_IdempotencyKey(t *testing.T) { }) } +func Test_sanitizeExtensionName(t *testing.T) { + tests := []struct { + name string + in string + want string + }{ + {name: "snake_case", in: "chain_id", want: "chainid"}, + {name: "dotted", in: "k8s.pod.name", want: "k8spodname"}, + {name: "already valid", in: "chainid", want: "chainid"}, + {name: "upper case is lowered", in: "ChainID", want: "chainid"}, + {name: "empty", in: "", want: ""}, + {name: "all invalid characters", in: "---...", want: ""}, + {name: "mixed valid and invalid", in: "Service-Name.1", want: "servicename1"}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, tt.want, sanitizeExtensionName(tt.in)) + }) + } +} + +func TestNewEventWithOpts_WithResourceAttributeExtensions(t *testing.T) { + payload := []byte("body") + + t.Run("sanitized keys/values land on the event", func(t *testing.T) { + attrs := map[string]string{"chain_id": "1", "k8s.pod.name": "pod-abc"} + event, err := NewEventWithOpts("domain", "entity", payload, nil, WithResourceAttributeExtensions(attrs)) + require.NoError(t, err) + ext := event.Extensions() + assert.Equal(t, "1", ext["chainid"]) + assert.Equal(t, "pod-abc", ext["k8spodname"]) + }) + + t.Run("empty sanitized name is dropped", func(t *testing.T) { + attrs := map[string]string{"---": "value"} + event, err := NewEventWithOpts("domain", "entity", payload, nil, WithResourceAttributeExtensions(attrs)) + require.NoError(t, err) + assert.Len(t, event.Extensions(), 1) // only the always-set recordedtime extension + }) + + t.Run("reserved name is skipped", func(t *testing.T) { + attrs := map[string]string{IdempotencyKeyAttr: "should-not-override", "subject": "should-not-override"} + event, err := NewEventWithOpts("domain", "entity", payload, map[string]any{IdempotencyKeyAttr: "real-key"}, WithResourceAttributeExtensions(attrs)) + require.NoError(t, err) + ext := event.Extensions() + assert.Equal(t, "real-key", ext[IdempotencyKeyAttr]) + assert.Empty(t, event.Subject()) + }) + + t.Run("duplicate sanitized names resolve deterministically to sorted-first key", func(t *testing.T) { + attrs := map[string]string{"service.name": "from-dotted", "service_name": "from-snake"} + event, err := NewEventWithOpts("domain", "entity", payload, nil, WithResourceAttributeExtensions(attrs)) + require.NoError(t, err) + // sorted order: "service.name" < "service_name" ('.' < '_' in ASCII), so the dotted key wins. + assert.Equal(t, "from-dotted", event.Extensions()["servicename"]) + }) + + t.Run("omitting all opts is a no-op", func(t *testing.T) { + event, err := NewEventWithOpts("domain", "entity", payload, nil) + require.NoError(t, err) + assert.Len(t, event.Extensions(), 1) // only the always-set recordedtime extension + }) +} + +// TestNewEvent_UnchangedSignature is a backward-compatibility guard: NewEvent's exported +// signature must stay exactly as it was before EventOpt/NewEventWithOpts were introduced, and +// must remain equivalent to calling NewEventWithOpts with no opts. +func TestNewEvent_UnchangedSignature(t *testing.T) { + payload := []byte("body") + attributes := map[string]any{"subject": "example-subject"} + + viaNewEvent, err := NewEvent("domain", "entity", payload, attributes) + require.NoError(t, err) + + viaNewEventWithOpts, err := NewEventWithOpts("domain", "entity", payload, attributes) + require.NoError(t, err) + + assert.Equal(t, viaNewEventWithOpts.Subject(), viaNewEvent.Subject()) + assert.Equal(t, viaNewEventWithOpts.Extensions()["recordedtime"].(ce.Timestamp).Truncate(time.Second), + viaNewEvent.Extensions()["recordedtime"].(ce.Timestamp).Truncate(time.Second)) + assert.Equal(t, viaNewEventWithOpts.Data(), viaNewEvent.Data()) +} + func TestEventToProto(t *testing.T) { // Create a test protobuf message testProto := pb.PingResponse{Message: "test message"} @@ -597,6 +681,19 @@ func TestOptions(t *testing.T) { assert.Equal(t, mockProvider, config.headerProvider) }) + t.Run("WithResourceAttributeHeaders", func(t *testing.T) { + config := defaultCfg + WithResourceAttributeHeaders(map[string]string{ + "Chain-ID": "1", + "id": "skipped", // reserved extension name + "chain_id": "2", // duplicate sanitized key, first wins + })(&config) + assert.NotNil(t, config.headerProvider) + headers, err := config.headerProvider.Headers(t.Context()) + require.NoError(t, err) + assert.Equal(t, map[string]string{"chainid": "1"}, headers) + }) + t.Run("WithBasicAuth", func(t *testing.T) { config := defaultCfg WithBasicAuth("user", "pass")(&config) @@ -668,6 +765,49 @@ func (m *mockHeaderProvider) Headers(ctx context.Context) (map[string]string, er return m.headers, nil } +// capturingServer is a minimal ChipIngressServer that records the incoming gRPC metadata +// of the last request it handles. +type capturingServer struct { + pb.UnimplementedChipIngressServer + lastMD metadata.MD +} + +func (s *capturingServer) Ping(ctx context.Context, _ *pb.EmptyRequest) (*pb.PingResponse, error) { + s.lastMD, _ = metadata.FromIncomingContext(ctx) + return &pb.PingResponse{}, nil +} + +// TestClient_ChainedHeaderProviders is a regression test for the fix that switched from +// grpc.WithUnaryInterceptor (last-one-wins) to grpc.WithChainUnaryInterceptor: when both +// WithHeaderProvider and WithNOPLookup are configured, both providers' headers must reach +// the server, not just the one registered last. +func TestClient_ChainedHeaderProviders(t *testing.T) { + lis, err := (&net.ListenConfig{}).Listen(t.Context(), "tcp", "127.0.0.1:0") + require.NoError(t, err) + defer lis.Close() + + srv := gp.NewServer() + capture := &capturingServer{} + pb.RegisterChipIngressServer(srv, capture) + go func() { _ = srv.Serve(lis) }() + defer srv.Stop() + + client, err := NewClient(lis.Addr().String(), + WithInsecureConnection(), + WithHeaderProvider(&mockHeaderProvider{headers: map[string]string{"x-resource-attr": "chain-1"}}), + WithNOPLookup(), + ) + require.NoError(t, err) + defer client.Close() //nolint:errcheck + + _, err = client.Ping(t.Context(), &EmptyRequest{}) + require.NoError(t, err) + + require.NotNil(t, capture.lastMD) + assert.Equal(t, []string{"chain-1"}, capture.lastMD.Get("x-resource-attr")) + assert.Equal(t, []string{"true"}, capture.lastMD.Get("x-include-nop-info")) +} + func TestWithTLS(t *testing.T) { serverName := "example.com" config := defaultCfg diff --git a/pkg/chipingress/header_provider.go b/pkg/chipingress/header_provider.go index 9a4141f98f..47e2dfcd9b 100644 --- a/pkg/chipingress/header_provider.go +++ b/pkg/chipingress/header_provider.go @@ -108,6 +108,55 @@ func newStaticHeaderProvider(headers map[string]string, requireTLS bool) HeaderP return &staticHeaderProvider{headers: headers, requireTLS: requireTLS} } +// NewStaticHeaderProvider returns a HeaderProvider that always returns the given headers, +// for use with WithHeaderProvider to attach fixed, non-auth gRPC metadata (e.g. resource +// attributes) to every request. +func NewStaticHeaderProvider(headers map[string]string) HeaderProvider { + return newStaticHeaderProvider(headers, false) +} + +// SanitizeMetadataValue replaces any byte outside the printable ASCII range [0x20-0x7E] +// with '?'. grpc-go hard-fails the entire RPC when an outgoing metadata value fails this +// check (unlike the CE-extension path, where an invalid entry is simply dropped), so +// values headed for gRPC metadata must be normalized before being sent. +func SanitizeMetadataValue(val string) string { + b := []byte(val) + out := make([]byte, len(b)) + for i, c := range b { + if c >= 0x20 && c <= 0x7E { + out[i] = c + } else { + out[i] = '?' + } + } + return string(out) +} + +// SanitizeMetadataHeaders sanitizes a map of resource-attribute headers for use as outgoing +// gRPC metadata (e.g. via NewStaticHeaderProvider). Keys are sanitized with +// sanitizeExtensionName — the same strict [a-z0-9] charset used for CloudEvent extensions — +// which is a subset of grpc's allowed metadata-key charset, so a sanitized key can never trip +// grpc's key validation or the reserved "-bin" suffix, and produces the same key stem as the +// corresponding CE extension (differing only by the CloudEvents Kafka binding's "ce_" prefix +// once on the wire). Values are sanitized via SanitizeMetadataValue, since grpc-go fails the +// whole RPC on a non-printable value. Entries that sanitize to an empty key, or that collide +// with a reserved extension name (see reservedExtensionNames) or a gRPC-reserved header name +// (see reservedMetadataKeys), are skipped. Keys are applied in sorted order so duplicate +// sanitized keys resolve deterministically (first in sorted order wins), matching +// WithResourceAttributeExtensions' collision handling. +// +// Note: unlike the CloudEvents Kafka binding, gRPC metadata keys are NOT prefixed with "ce_" — +// that prefix is a CloudEvents-binding concept, not a metadata one, and reusing it here would +// collide with the CE binding's own "ce_" Kafka header if the server ever forwards gRPC +// metadata verbatim onto Kafka. +func SanitizeMetadataHeaders(in map[string]string) map[string]string { + out := make(map[string]string, len(in)) + for _, pair := range sanitizeResourceAttributeKeys(in, reservedMetadataKeys) { + out[pair.name] = SanitizeMetadataValue(in[pair.key]) + } + return out +} + // newRotatingHeaderProvider returns a HeaderProvider that refreshes its // headers every ttl using signer. initialHeaders, if non-empty, are served // until the first rotation occurs. diff --git a/pkg/chipingress/header_provider_test.go b/pkg/chipingress/header_provider_test.go index 0c0f523e90..8069420fd2 100644 --- a/pkg/chipingress/header_provider_test.go +++ b/pkg/chipingress/header_provider_test.go @@ -4,13 +4,18 @@ import ( "context" "crypto/ed25519" "encoding/hex" + "net" "testing" "time" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "google.golang.org/grpc" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" "github.com/smartcontractkit/chainlink-common/pkg/chipingress" + "github.com/smartcontractkit/chainlink-common/pkg/chipingress/pb" ) // fakeSigner is a minimal chipingress.Signer used by tests that need to @@ -260,3 +265,132 @@ func TestNewHeaderProvider(t *testing.T) { assert.Nil(t, provider) }) } + +func TestNewStaticHeaderProvider(t *testing.T) { + headers := map[string]string{"chain_id": "1", "environment": "prod"} + provider := chipingress.NewStaticHeaderProvider(headers) + require.NotNil(t, provider) + + got, err := provider.Headers(t.Context()) + require.NoError(t, err) + assert.Equal(t, headers, got) + + type tsr interface { + RequireTransportSecurity() bool + } + tlsReq, ok := provider.(tsr) + require.True(t, ok) + assert.False(t, tlsReq.RequireTransportSecurity()) +} + +func TestSanitizeMetadataValue(t *testing.T) { + tests := []struct { + name string + in string + want string + }{ + {name: "printable ASCII is unchanged", in: "chain-1_prod.v2", want: "chain-1_prod.v2"}, + {name: "empty", in: "", want: ""}, + {name: "control character replaced", in: "value\nwith\tcontrol", want: "value?with?control"}, + {name: "non-ASCII UTF-8 replaced byte-wise", in: "café", want: "caf??"}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, tt.want, chipingress.SanitizeMetadataValue(tt.in)) + }) + } +} + +func TestSanitizeMetadataHeaders(t *testing.T) { + t.Run("standard OTel-style keys are sanitized to the same stem as CE extensions", func(t *testing.T) { + in := map[string]string{ + "service.name": "beholder", + "chain_id": "1", + "node-operator": "acme", + } + got := chipingress.SanitizeMetadataHeaders(in) + assert.Equal(t, map[string]string{ + "servicename": "beholder", + "chainid": "1", + "nodeoperator": "acme", + }, got) + }) + + t.Run("empty-after-sanitize keys are dropped", func(t *testing.T) { + got := chipingress.SanitizeMetadataHeaders(map[string]string{"---": "value"}) + assert.Empty(t, got) + }) + + t.Run("reserved names are dropped", func(t *testing.T) { + got := chipingress.SanitizeMetadataHeaders(map[string]string{chipingress.IdempotencyKeyAttr: "should-not-appear", "subject": "should-not-appear"}) + assert.Empty(t, got) + }) + + t.Run("gRPC-reserved header 'te' is dropped", func(t *testing.T) { + got := chipingress.SanitizeMetadataHeaders(map[string]string{"te": "trailers"}) + assert.Empty(t, got) + }) + + t.Run("non-printable values are sanitized", func(t *testing.T) { + got := chipingress.SanitizeMetadataHeaders(map[string]string{"chain_id": "1\n2"}) + assert.Equal(t, "1?2", got["chainid"]) + }) + + t.Run("duplicate sanitized keys resolve deterministically to sorted-first key", func(t *testing.T) { + got := chipingress.SanitizeMetadataHeaders(map[string]string{"service.name": "from-dotted", "service_name": "from-snake"}) + // sorted order: "service.name" < "service_name" ('.' < '_' in ASCII), so the dotted key wins. + assert.Equal(t, "from-dotted", got["servicename"]) + }) +} + +// pingServer is a minimal ChipIngressServer that always answers Ping successfully. +type pingServer struct { + pb.UnimplementedChipIngressServer +} + +func (pingServer) Ping(context.Context, *pb.EmptyRequest) (*pb.PingResponse, error) { + return &pb.PingResponse{}, nil +} + +// TestSanitizeMetadataHeaders_AvoidsRPCFailure is a regression/guard test for the core reason +// SanitizeMetadataHeaders exists: grpc-go hard-fails an entire RPC (codes.Internal) when an +// outgoing metadata pair fails its charset validation. An unsanitized resource-attribute key or +// value (dots, non-printable characters) reproduces that failure; running it through +// SanitizeMetadataHeaders first must not. +func TestSanitizeMetadataHeaders_AvoidsRPCFailure(t *testing.T) { + lis, err := (&net.ListenConfig{}).Listen(t.Context(), "tcp", "127.0.0.1:0") + require.NoError(t, err) + defer lis.Close() + + srv := grpc.NewServer() + pb.RegisterChipIngressServer(srv, pingServer{}) + go func() { _ = srv.Serve(lis) }() + defer srv.Stop() + + dirty := map[string]string{"k8s.pod.name": "pod-\x01abc"} + + t.Run("unsanitized headers fail the RPC", func(t *testing.T) { + client, err := chipingress.NewClient(lis.Addr().String(), + chipingress.WithInsecureConnection(), + chipingress.WithHeaderProvider(chipingress.NewStaticHeaderProvider(dirty)), + ) + require.NoError(t, err) + defer client.Close() //nolint:errcheck + + _, err = client.Ping(t.Context(), &chipingress.EmptyRequest{}) + require.Error(t, err) + assert.Equal(t, codes.Internal, status.Code(err)) + }) + + t.Run("sanitized headers succeed", func(t *testing.T) { + client, err := chipingress.NewClient(lis.Addr().String(), + chipingress.WithInsecureConnection(), + chipingress.WithHeaderProvider(chipingress.NewStaticHeaderProvider(chipingress.SanitizeMetadataHeaders(dirty))), + ) + require.NoError(t, err) + defer client.Close() //nolint:errcheck + + _, err = client.Ping(t.Context(), &chipingress.EmptyRequest{}) + require.NoError(t, err) + }) +} diff --git a/pkg/chipingress/resource_attributes.go b/pkg/chipingress/resource_attributes.go new file mode 100644 index 0000000000..2d5f974e8f --- /dev/null +++ b/pkg/chipingress/resource_attributes.go @@ -0,0 +1,46 @@ +package chipingress + +import "sort" + +// resourceAttrKey pairs a sanitized extension/metadata key name with the original +// resource-attribute key it was derived from. +type resourceAttrKey struct { + name string + key string +} + +// sanitizeResourceAttributeKeys returns the deduplicated, sorted list of resource-attribute +// keys that survive sanitization and reservation checks. The returned pairs contain the +// sanitized name and the original map key, so callers can apply their own value handling. +// +// Ordering is deterministic: original keys are sorted lexicographically, and if two keys +// sanitize to the same name the first one in sorted order wins. extraReserved, if non-nil, +// is consulted in addition to reservedExtensionNames. +func sanitizeResourceAttributeKeys(attrs map[string]string, extraReserved map[string]struct{}) []resourceAttrKey { + keys := make([]string, 0, len(attrs)) + for k := range attrs { + keys = append(keys, k) + } + sort.Strings(keys) + + seen := make(map[string]struct{}, len(attrs)) + result := make([]resourceAttrKey, 0, len(attrs)) + for _, k := range keys { + name := sanitizeExtensionName(k) + if name == "" { + continue + } + if _, reserved := reservedExtensionNames[name]; reserved { + continue + } + if _, reserved := extraReserved[name]; reserved { + continue + } + if _, already := seen[name]; already { + continue + } + seen[name] = struct{}{} + result = append(result, resourceAttrKey{name: name, key: k}) + } + return result +} diff --git a/pkg/chipingress/types.go b/pkg/chipingress/types.go index e27c18d194..25eede3edb 100644 --- a/pkg/chipingress/types.go +++ b/pkg/chipingress/types.go @@ -13,6 +13,35 @@ import ( // Kafka headers named "ce_" (e.g., ce_idempotencykey), enabling downstream deduplication. const IdempotencyKeyAttr = "idempotencykey" +// reservedExtensionNames holds every CloudEvent extension name that NewEvent sets internally, +// plus the CloudEvents core context attribute names (id, source, type, specversion, time, +// subject, dataschema, datacontenttype) and the spec-forbidden "data" name. WithResourceAttributeExtensions +// consults this set so that a resource attribute can never silently overwrite event-lifecycle +// metadata or collide with a CloudEvents core attribute. +var reservedExtensionNames = map[string]struct{}{ + IdempotencyKeyAttr: {}, + "recordedtime": {}, + "id": {}, + "source": {}, + "type": {}, + "specversion": {}, + "time": {}, + "subject": {}, + "dataschema": {}, + "datacontenttype": {}, + "data": {}, +} + +// reservedMetadataKeys holds gRPC-reserved header names that could otherwise be reached by +// sanitizeExtensionName's [a-z0-9] sanitization. Verified against grpc-go v1.79.1's +// isReservedHeader: every other reserved header (pseudo-headers, "content-type", "grpc-*") +// contains a ':' or '-' that sanitization strips, so "te" is the only one actually reachable. +// SanitizeMetadataHeaders consults this set so that edge case is handled deterministically +// rather than relying on grpc's own (silent) handling of a reserved header. +var reservedMetadataKeys = map[string]struct{}{ + "te": {}, +} + type ( // Cloudevents types CloudEvent = ce.Event