Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
59 changes: 16 additions & 43 deletions pkg/chipingress/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@ import (
"crypto/tls"
"fmt"
"net"
"strings"
"time"

"github.com/google/uuid"
Expand Down Expand Up @@ -215,9 +214,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.
// WithResourceAttributeHeaders returns an Opt that attaches the provided resource attributes as
// gRPC metadata on every request, under ResourceHeaderPrefix. It combines SanitizeMetadataHeaders
// with NewStaticHeaderProvider so the safe, validated path is used by default.
//
// Attributes are attached once per request rather than to individual events because they describe the
// producer, not any one event. Chip-ingress fans them out onto every Kafka record the request
// produces.
func WithResourceAttributeHeaders(attrs map[string]string) Opt {
return WithHeaderProvider(NewStaticHeaderProvider(SanitizeMetadataHeaders(attrs)))
}
Expand Down Expand Up @@ -254,11 +257,15 @@ func WithTracerProvider(provider trace.TracerProvider) Opt {
return func(c *clientConfig) { c.tracerProvider = provider }
}

// nopInfoHeaderKey is the metadata key WithNOPLookup sets, asking chip-ingress to look up NOP info
// for the authenticated CSA key.
const nopInfoHeaderKey = "x-include-nop-info"

func WithNOPLookup() Opt {
return func(c *clientConfig) {
c.nopInfoHeaderProvider = headerProviderFunc(func(ctx context.Context) (map[string]string, error) {
return map[string]string{
"x-include-nop-info": "true",
nopInfoHeaderKey: "true",
}, nil
})
}
Expand All @@ -283,42 +290,12 @@ 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.
//
// Resource attributes are deliberately not stamped here. They describe the producer rather than any
// individual event, so they travel once per request as gRPC metadata (see
// WithResourceAttributeHeaders) instead of being repeated on every event in a batch.
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)
Expand Down Expand Up @@ -352,10 +329,6 @@ func NewEventWithOpts(domain, entity string, payload []byte, attributes map[stri
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)
Expand Down
149 changes: 62 additions & 87 deletions pkg/chipingress/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"context"
"fmt"
"net"
"strings"
"testing"
"time"

Expand Down Expand Up @@ -168,89 +169,6 @@ 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"}
Expand Down Expand Up @@ -684,14 +602,21 @@ func TestOptions(t *testing.T) {
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
"Chain-ID": "1", // lower-cased, separator preserved
"csa_public_key": "abc", // preserved verbatim
// Namespaced rather than dropped: prefixing puts them out of reach of the real keys.
"te": "harmless",
authHeaderKey: "harmless",
})(&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)
assert.Equal(t, map[string]string{
ResourceHeaderPrefix + "chain-id": "1",
ResourceHeaderPrefix + "csa_public_key": "abc",
ResourceHeaderPrefix + "te": "harmless",
ResourceHeaderPrefix + strings.ToLower(authHeaderKey): "harmless",
}, headers)
})

t.Run("WithBasicAuth", func(t *testing.T) {
Expand Down Expand Up @@ -808,6 +733,56 @@ func TestClient_ChainedHeaderProviders(t *testing.T) {
assert.Equal(t, []string{"true"}, capture.lastMD.Get("x-include-nop-info"))
}

// TestClient_AuthHeaderCoexistsWithResourceAttributes pins down the property the resource-attribute
// work must never break: the CSA node auth token and the resource-attribute headers travel by two
// different mechanisms — per-RPC credentials (WithTokenAuth) and a unary interceptor
// (WithResourceAttributeHeaders) — and both must arrive intact, exactly once, on the same request.
//
// It also pins the property that lets the client carry attributes without a reserved-key deny-list:
// an attribute named after the auth header is namespaced under ResourceHeaderPrefix, so it cannot
// append a second value under the auth header's own key.
func TestClient_AuthHeaderCoexistsWithResourceAttributes(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()

const authToken = "2:deadbeef:1:cafe"

client, err := NewClient(lis.Addr().String(),
WithInsecureConnection(),
WithTokenAuth(&mockHeaderProvider{headers: map[string]string{authHeaderKey: authToken}}),
WithResourceAttributeHeaders(map[string]string{
"csa_public_key": "abc123",
"service.name": "chainlink",
// Namespaced away from the auth key rather than appended to it.
authHeaderKey: "forged",
}),
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)
// grpc lower-cases metadata keys on the wire.
assert.Equal(t, []string{authToken}, capture.lastMD.Get(authHeaderKey),
"the auth token must arrive exactly once, unmodified")
assert.Equal(t, []string{"abc123"}, capture.lastMD.Get(ResourceHeaderPrefix+"csa_public_key"))
assert.Equal(t, []string{"chainlink"}, capture.lastMD.Get(ResourceHeaderPrefix+"service.name"))
assert.Equal(t, []string{"true"}, capture.lastMD.Get("x-include-nop-info"))
// The forged attribute landed in the resource namespace, harmlessly.
assert.Equal(t, []string{"forged"},
capture.lastMD.Get(ResourceHeaderPrefix+strings.ToLower(authHeaderKey)))
}

func TestWithTLS(t *testing.T) {
serverName := "example.com"
config := defaultCfg
Expand Down
Loading
Loading