diff --git a/protocol/pubsub/v2/doc.go b/protocol/pubsub/v2/doc.go index db8467b0c..a4936312d 100644 --- a/protocol/pubsub/v2/doc.go +++ b/protocol/pubsub/v2/doc.go @@ -8,5 +8,23 @@ Package pubsub implements a Pub/Sub binding using google.cloud.com/go/pubsub mod PubSub Messages can be modified beyond what CloudEvents cover by using `WithOrderingKey` or `WithCustomAttributes`. See function docs for more details. + +In binary mode, content-type and ce-datacontenttype both describe the event data's +media type. If both are present, their values must match exactly; ReadBinary +rejects conflicting values. When reading, legacy Content-Type is used only when +neither is present. Binary messages are written with ce-datacontenttype and, for +backward compatibility with existing consumers, Content-Type. Use +ce-datacontenttype when filtering binary-mode events by data content type. + +For structured JSON, the Google Cloud Pub/Sub binding defines the content-type +attribute as the CloudEvents envelope's media type (for example, +application/cloudevents+json; charset=utf-8). The datacontenttype property inside +the JSON describes the event data's media type. The binding permits this property +to also be included as a ce-datacontenttype Pub/Sub attribute, but that duplication +is optional in structured mode and cannot be assumed when defining filters. + +When reading messages, content-type determines the encoding. If it is absent, +legacy Content-Type is used instead. A known CloudEvents event format selects +structured mode even when ce-specversion is also present. */ package pubsub diff --git a/protocol/pubsub/v2/message.go b/protocol/pubsub/v2/message.go index fa4ffe4fb..1d32a950a 100644 --- a/protocol/pubsub/v2/message.go +++ b/protocol/pubsub/v2/message.go @@ -8,6 +8,7 @@ package pubsub import ( "bytes" "context" + "fmt" "strings" "cloud.google.com/go/pubsub/v2" @@ -18,8 +19,12 @@ import ( ) const ( - prefix = "ce-" - contentType = "Content-Type" + prefix = "ce-" + + // contentType identifies the event format in structured mode and the data media type in binary mode. + contentType = "content-type" + // legacyContentType is retained for compatibility with older senders and consumers. + legacyContentType = "Content-Type" ) var specs = spec.WithPrefix(prefix) @@ -38,12 +43,14 @@ func NewMessage(pm *pubsub.Message) *Message { var f format.Format = nil var version spec.Version = nil if pm.Attributes != nil { - // Use Content-type attr to determine if message is structured and - // set format. - if s := pm.Attributes[contentType]; format.IsFormat(s) { + s, ok := pm.Attributes[contentType] + if !ok { + // Use the legacy key only when the binding's content-type is absent. + s = pm.Attributes[legacyContentType] + } + if format.IsFormat(s) { f = format.Lookup(s) } - // Binary v0.3: if s := pm.Attributes[specs.PrefixedSpecVersionName()]; s != "" { version = specs.Version(s) } @@ -61,29 +68,30 @@ var _ binding.Message = (*Message)(nil) var _ binding.MessageMetadataReader = (*Message)(nil) func (m *Message) ReadEncoding() binding.Encoding { - if m.version != nil { - return binding.EncodingBinary - } + // Structured messages may duplicate CloudEvents attributes in metadata. if m.format != nil { return binding.EncodingStructured } + if m.version != nil { + return binding.EncodingBinary + } return binding.EncodingUnknown } func (m *Message) ReadStructured(ctx context.Context, encoder binding.StructuredWriter) error { - if m.version != nil { - return binding.ErrNotStructured - } - if m.format == nil { + if m.ReadEncoding() != binding.EncodingStructured { return binding.ErrNotStructured } return encoder.SetStructuredEvent(ctx, m.format, bytes.NewReader(m.internal.Data)) } func (m *Message) ReadBinary(ctx context.Context, encoder binding.BinaryWriter) (err error) { - if m.format != nil { + if m.ReadEncoding() != binding.EncodingBinary { return binding.ErrNotBinary } + if _, err = m.binaryDataContentType(m.version.AttributeFromKind(spec.DataContentType)); err != nil { + return err + } for k, v := range m.internal.Attributes { if strings.HasPrefix(k, prefix) { @@ -93,8 +101,19 @@ func (m *Message) ReadBinary(ctx context.Context, encoder binding.BinaryWriter) } else { err = encoder.SetExtension(strings.TrimPrefix(k, prefix), string(v)) } - } else if k == contentType { - err = encoder.SetAttribute(m.version.AttributeFromKind(spec.DataContentType), string(v)) + } else if k == contentType || k == legacyContentType { + attr := m.version.AttributeFromKind(spec.DataContentType) + // Let the prefixed attribute provide the value when present. + // content-type was checked for agreement; Content-Type is only a fallback. + if _, ok := m.internal.Attributes[prefix+attr.Name()]; ok { + continue + } + if k == legacyContentType { + if _, ok := m.internal.Attributes[contentType]; ok { + continue + } + } + err = encoder.SetAttribute(attr, string(v)) } if err != nil { return err @@ -111,11 +130,35 @@ func (m *Message) ReadBinary(ctx context.Context, encoder binding.BinaryWriter) func (m *Message) GetAttribute(k spec.Kind) (spec.Attribute, interface{}) { attr := m.version.AttributeFromKind(k) if attr != nil { + if k == spec.DataContentType { + value, err := m.binaryDataContentType(attr) + if err != nil { + // There is no unambiguous value; ReadBinary reports the conflict. + return attr, nil + } + return attr, value + } return attr, m.internal.Attributes[prefix+attr.Name()] } return nil, nil } +func (m *Message) binaryDataContentType(attr spec.Attribute) (string, error) { + prefixedName := prefix + attr.Name() + prefixedValue, hasPrefixedValue := m.internal.Attributes[prefixedName] + messageValue, hasMessageValue := m.internal.Attributes[contentType] + if hasPrefixedValue && hasMessageValue && prefixedValue != messageValue { + return "", fmt.Errorf("pubsub: conflicting %q and %q values: %q != %q", contentType, prefixedName, messageValue, prefixedValue) + } + if hasMessageValue { + return messageValue, nil + } + if hasPrefixedValue { + return prefixedValue, nil + } + return m.internal.Attributes[legacyContentType], nil +} + func (m *Message) GetExtension(name string) interface{} { return m.internal.Attributes[prefix+name] } diff --git a/protocol/pubsub/v2/message_test.go b/protocol/pubsub/v2/message_test.go index 22821c370..950dc54dc 100644 --- a/protocol/pubsub/v2/message_test.go +++ b/protocol/pubsub/v2/message_test.go @@ -12,10 +12,411 @@ import ( "cloud.google.com/go/pubsub/v2" "github.com/cloudevents/sdk-go/v2/binding" + "github.com/cloudevents/sdk-go/v2/binding/spec" "github.com/cloudevents/sdk-go/v2/event" "github.com/cloudevents/sdk-go/v2/protocol" + "github.com/cloudevents/sdk-go/v2/test" + "github.com/stretchr/testify/require" ) +func TestReadBinaryDataContentType(t *testing.T) { + tests := []struct { + name string + attributes map[string]string + want string + }{ + { + name: "prefixed attribute", + attributes: map[string]string{"ce-datacontenttype": "application/json; charset=utf-8"}, + want: "application/json; charset=utf-8", + }, + { + name: "legacy attribute", + attributes: map[string]string{"Content-Type": "application/json; charset=utf-8"}, + want: "application/json; charset=utf-8", + }, + { + name: "lowercase content type", + attributes: map[string]string{"content-type": "application/json; charset=utf-8"}, + want: "application/json; charset=utf-8", + }, + { + name: "lowercase content type takes precedence over legacy attribute", + attributes: map[string]string{ + "content-type": "application/json; charset=utf-8", + "Content-Type": "text/plain", + }, + want: "application/json; charset=utf-8", + }, + { + name: "matching lowercase and prefixed attributes", + attributes: map[string]string{ + "ce-datacontenttype": "application/json; charset=utf-8", + "content-type": "application/json; charset=utf-8", + }, + want: "application/json; charset=utf-8", + }, + { + name: "matching lowercase and prefixed attributes override legacy attribute", + attributes: map[string]string{ + "ce-datacontenttype": "application/json; charset=utf-8", + "content-type": "application/json; charset=utf-8", + "Content-Type": "text/plain", + }, + want: "application/json; charset=utf-8", + }, + { + name: "all content type attributes match", + attributes: map[string]string{ + "ce-datacontenttype": "application/json; charset=utf-8", + "content-type": "application/json; charset=utf-8", + "Content-Type": "application/json; charset=utf-8", + }, + want: "application/json; charset=utf-8", + }, + { + name: "matching empty lowercase and prefixed attributes", + attributes: map[string]string{ + "ce-datacontenttype": "", + "content-type": "", + }, + }, + { + name: "empty lowercase content type takes precedence over legacy attribute", + attributes: map[string]string{ + "content-type": "", + "Content-Type": "text/plain", + }, + }, + { + name: "matching prefixed and legacy attributes", + attributes: map[string]string{ + "ce-datacontenttype": "application/json; charset=utf-8", + "Content-Type": "application/json; charset=utf-8", + }, + want: "application/json; charset=utf-8", + }, + { + name: "prefixed attribute takes precedence over legacy attribute", + attributes: map[string]string{ + "ce-datacontenttype": "application/json; charset=utf-8", + "Content-Type": "text/plain", + }, + want: "application/json; charset=utf-8", + }, + { + name: "empty prefixed attribute takes precedence over legacy attribute", + attributes: map[string]string{ + "ce-datacontenttype": "", + "Content-Type": "text/plain", + }, + }, + { + name: "empty legacy attribute", + attributes: map[string]string{"Content-Type": ""}, + }, + {name: "no content type"}, + } + for _, version := range []string{event.CloudEventsVersionV03, event.CloudEventsVersionV1} { + t.Run(version, func(t *testing.T) { + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + attributes := map[string]string{ + "ce-specversion": version, + "ce-id": "testid", + "ce-source": "/test", + "ce-type": "test.type", + } + for name, value := range tt.attributes { + attributes[name] = value + } + data := []byte(`{"hello":"world"}`) + msg := NewMessage(&pubsub.Message{ + Attributes: attributes, + Data: data, + }) + + require.Equal(t, binding.EncodingBinary, msg.ReadEncoding()) + attr, value := msg.GetAttribute(spec.DataContentType) + require.Equal(t, spec.DataContentType, attr.Kind()) + require.Equal(t, tt.want, value) + got, err := binding.ToEvent(context.Background(), msg) + require.NoError(t, err) + require.Equal(t, tt.want, got.DataContentType()) + require.Equal(t, data, got.Data()) + }) + } + }) + } +} + +func TestReadBinaryConflictingDataContentTypes(t *testing.T) { + tests := []struct { + name string + attributes map[string]string + }{ + { + name: "conflicting lowercase and prefixed attributes", + attributes: map[string]string{ + "content-type": "application/json", + "ce-datacontenttype": "text/plain", + }, + }, + { + name: "empty lowercase attribute conflicts with prefixed attribute", + attributes: map[string]string{ + "content-type": "", + "ce-datacontenttype": "text/plain", + }, + }, + { + name: "empty prefixed attribute conflicts with lowercase attribute", + attributes: map[string]string{ + "content-type": "application/json", + "ce-datacontenttype": "", + }, + }, + { + name: "legacy attribute does not resolve a conflict", + attributes: map[string]string{ + "content-type": "application/json", + "ce-datacontenttype": "text/plain", + "Content-Type": "application/json", + }, + }, + } + for _, version := range []string{event.CloudEventsVersionV03, event.CloudEventsVersionV1} { + t.Run(version, func(t *testing.T) { + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + attributes := map[string]string{ + "ce-specversion": version, + "ce-id": "testid", + "ce-source": "/test", + "ce-type": "test.type", + } + for name, value := range tt.attributes { + attributes[name] = value + } + msg := NewMessage(&pubsub.Message{Attributes: attributes, Data: []byte(`{"hello":"world"}`)}) + require.Equal(t, binding.EncodingBinary, msg.ReadEncoding()) + + writer := &pubsubMessagePublisher{Attributes: make(map[string]string)} + err := msg.ReadBinary(context.Background(), writer) + require.ErrorContains(t, err, `conflicting "content-type" and "ce-datacontenttype"`) + require.Empty(t, writer.Attributes) + require.Nil(t, writer.Data) + attr, value := msg.GetAttribute(spec.DataContentType) + require.Equal(t, spec.DataContentType, attr.Kind()) + require.Nil(t, value) + + got, conversionErr := binding.ToEvent(context.Background(), msg) + require.EqualError(t, conversionErr, err.Error()) + require.Nil(t, got) + for _, encoding := range []binding.Encoding{binding.EncodingBinary, binding.EncodingStructured} { + t.Run(encoding.String(), func(t *testing.T) { + ctx := binding.WithForceBinary(context.Background()) + if encoding == binding.EncodingStructured { + ctx = binding.WithForceStructured(context.Background()) + } + pm := &pubsub.Message{Attributes: make(map[string]string)} + require.EqualError(t, WritePubSubMessage(ctx, msg, pm), err.Error()) + require.Empty(t, pm.Attributes) + require.Nil(t, pm.Data) + }) + } + }) + } + }) + } +} + +func TestNewMessageContentType(t *testing.T) { + tests := []struct { + name string + attributes map[string]string + encoding binding.Encoding + }{ + { + name: "lowercase structured content type", + attributes: map[string]string{"content-type": "application/cloudevents+json; charset=utf-8"}, + encoding: binding.EncodingStructured, + }, + { + name: "legacy structured content type", + attributes: map[string]string{"Content-Type": "application/cloudevents+json"}, + encoding: binding.EncodingStructured, + }, + { + name: "matching structured content types", + attributes: map[string]string{ + "content-type": "application/cloudevents+json", + "Content-Type": "application/cloudevents+json", + }, + encoding: binding.EncodingStructured, + }, + { + name: "lowercase structured content type takes precedence over legacy content type", + attributes: map[string]string{ + "content-type": "application/cloudevents+json", + "Content-Type": "application/json", + }, + encoding: binding.EncodingStructured, + }, + { + name: "lowercase structured content type with ce-specversion", + attributes: map[string]string{ + "content-type": "application/cloudevents+json", + "ce-specversion": "1.0", + }, + encoding: binding.EncodingStructured, + }, + { + name: "lowercase binary content type takes precedence over legacy content type", + attributes: map[string]string{ + "content-type": "application/json", + "Content-Type": "application/cloudevents+json", + "ce-specversion": "1.0", + }, + encoding: binding.EncodingBinary, + }, + { + name: "empty lowercase content type takes precedence over legacy content type", + attributes: map[string]string{ + "content-type": "", + "Content-Type": "application/cloudevents+json", + "ce-specversion": "1.0", + }, + encoding: binding.EncodingBinary, + }, + { + name: "legacy structured content type with ce-specversion", + attributes: map[string]string{ + "Content-Type": "application/cloudevents+json", + "ce-specversion": "1.0", + }, + encoding: binding.EncodingStructured, + }, + { + name: "unsupported lowercase format does not fall back to legacy format", + attributes: map[string]string{ + "content-type": "application/cloudevents+unknown", + "Content-Type": "application/cloudevents+json", + "ce-specversion": "1.0", + }, + encoding: binding.EncodingBinary, + }, + { + name: "lowercase content type without binary metadata", + attributes: map[string]string{ + "content-type": "application/json", + "Content-Type": "application/cloudevents+json", + }, + encoding: binding.EncodingUnknown, + }, + {name: "no attributes", encoding: binding.EncodingUnknown}, + { + name: "unsupported spec version", + attributes: map[string]string{"ce-specversion": "unknown"}, + encoding: binding.EncodingUnknown, + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + msg := NewMessage(&pubsub.Message{Attributes: tt.attributes}) + require.Equal(t, tt.encoding, msg.ReadEncoding()) + writer := &pubsubMessagePublisher{Attributes: make(map[string]string)} + err := msg.ReadStructured(context.Background(), writer) + if tt.encoding == binding.EncodingStructured { + require.NoError(t, err) + } else { + require.ErrorIs(t, err, binding.ErrNotStructured) + } + err = msg.ReadBinary(context.Background(), writer) + if tt.encoding == binding.EncodingBinary { + require.NoError(t, err) + } else { + require.ErrorIs(t, err, binding.ErrNotBinary) + } + }) + } +} + +func TestReadStructuredWithCloudEventAttributes(t *testing.T) { + tests := []struct { + name string + contentTypes map[string]string + }{ + { + name: "lowercase content type only", + contentTypes: map[string]string{ + "content-type": "application/cloudevents+json; charset=utf-8", + }, + }, + { + name: "legacy content type only", + contentTypes: map[string]string{ + "Content-Type": "application/cloudevents+json; charset=utf-8", + }, + }, + { + name: "matching lowercase and legacy content types", + contentTypes: map[string]string{ + "content-type": "application/cloudevents+json; charset=utf-8", + "Content-Type": "application/cloudevents+json; charset=utf-8", + }, + }, + { + name: "lowercase content type overrides legacy content type", + contentTypes: map[string]string{ + "content-type": "application/cloudevents+json; charset=utf-8", + "Content-Type": "application/protobuf", + }, + }, + } + test.EachEvent(t, test.Events(), func(t *testing.T, eventIn event.Event) { + eventIn = test.ConvertEventExtensionsToString(t, eventIn) + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + attributes := map[string]string{ + "ce-specversion": eventIn.SpecVersion(), + "ce-id": "ignored-id", + "ce-datacontenttype": "application/protobuf", + } + for name, value := range tt.contentTypes { + attributes[name] = value + } + msg := NewMessage(&pubsub.Message{ + Attributes: attributes, + Data: test.MustJSON(t, eventIn), + }) + require.Equal(t, binding.EncodingStructured, msg.ReadEncoding()) + eventOut, err := binding.ToEvent(context.Background(), msg) + require.NoError(t, err) + test.AssertEventEquals(t, eventIn, *eventOut) + }) + } + }) +} + +func TestReadStructuredMissingSpecVersion(t *testing.T) { + for _, name := range []string{"content-type", "Content-Type"} { + t.Run(name, func(t *testing.T) { + msg := NewMessage(&pubsub.Message{ + Attributes: map[string]string{ + name: event.ApplicationCloudEventsJSON, + "ce-specversion": event.CloudEventsVersionV1, + }, + Data: []byte(`{"id":"testid","source":"/test","type":"test.type"}`), + }) + require.Equal(t, binding.EncodingStructured, msg.ReadEncoding()) + eventOut, err := binding.ToEvent(context.Background(), msg) + require.ErrorContains(t, err, "no specversion") + require.Nil(t, eventOut) + }) + } +} + func TestReadStructured(t *testing.T) { tests := []struct { name string @@ -33,7 +434,7 @@ func TestReadStructured(t *testing.T) { name: "json format", pm: &pubsub.Message{ ID: "testid", - Attributes: map[string]string{contentType: event.ApplicationCloudEventsJSON}, + Attributes: map[string]string{"Content-Type": event.ApplicationCloudEventsJSON}, }, }, } diff --git a/protocol/pubsub/v2/protocol_test.go b/protocol/pubsub/v2/protocol_test.go index a65cd0a6a..094218090 100644 --- a/protocol/pubsub/v2/protocol_test.go +++ b/protocol/pubsub/v2/protocol_test.go @@ -26,6 +26,7 @@ func (pc *testPubsubClient) NewWithAttributesInterceptor(ctx context.Context, pr pc.srv = pstest.NewServer() conn, err := grpc.NewClient(pc.srv.Addr, grpc.WithTransportCredentials(insecure.NewCredentials()), grpc.WithUnaryInterceptor(customAttributesInterceptor(map[string]string{ "Content-Type": "text/json", + "ce-datacontenttype": "text/json", "ce-dataschema": "http://example.com/schema", "ce-exbinary": "AAECAw==", "ce-exbool": "true", diff --git a/protocol/pubsub/v2/write_pubsub_message.go b/protocol/pubsub/v2/write_pubsub_message.go index 85ab97df7..125ac0b89 100644 --- a/protocol/pubsub/v2/write_pubsub_message.go +++ b/protocol/pubsub/v2/write_pubsub_message.go @@ -9,6 +9,7 @@ import ( "bytes" "context" "io" + "strings" "cloud.google.com/go/pubsub/v2" "github.com/cloudevents/sdk-go/v2/binding" @@ -41,11 +42,19 @@ func (b *pubsubMessagePublisher) SetStructuredEvent(ctx context.Context, f forma if err != nil { return err } + if b.Attributes == nil { + b.Attributes = make(map[string]string) + } + b.clearCloudEventAttributes() + b.Attributes[contentType] = f.MediaType() b.Data = buf.Bytes() return nil } func (b *pubsubMessagePublisher) Start(ctx context.Context) error { + // Attributes and data omitted by the next event must not survive reuse. + b.clearCloudEventAttributes() + b.Data = nil return nil } @@ -53,6 +62,15 @@ func (b *pubsubMessagePublisher) End(ctx context.Context) error { return nil } +func (b *pubsubMessagePublisher) clearCloudEventAttributes() { + // Preserve custom Pub/Sub attributes while removing metadata from the previous event. + for name := range b.Attributes { + if strings.HasPrefix(name, prefix) || name == contentType || name == legacyContentType { + delete(b.Attributes, name) + } + } +} + func (b *pubsubMessagePublisher) SetData(reader io.Reader) error { buf, ok := reader.(*bytes.Buffer) if !ok { @@ -67,28 +85,24 @@ func (b *pubsubMessagePublisher) SetData(reader io.Reader) error { } func (b *pubsubMessagePublisher) SetAttribute(attribute spec.Attribute, value interface{}) error { - if attribute.Kind() == spec.DataContentType { - if value == nil { - delete(b.Attributes, contentType) - } - - // Everything is a string here - s, err := types.Format(value) - if err != nil { - return err - } - b.Attributes[contentType] = s - } else { - if value == nil { - delete(b.Attributes, prefix+attribute.Name()) + if value == nil { + delete(b.Attributes, prefix+attribute.Name()) + if attribute.Kind() == spec.DataContentType { + delete(b.Attributes, legacyContentType) } + return nil + } - // Everything is a string here - s, err := types.Format(value) - if err != nil { - return err - } - b.Attributes[prefix+attribute.Name()] = s + // Everything is a string here + s, err := types.Format(value) + if err != nil { + return err + } + b.Attributes[prefix+attribute.Name()] = s + if attribute.Kind() == spec.DataContentType { + // Retain Content-Type for backward compatibility with existing consumers. + // Use ce-datacontenttype when filtering binary-mode events by data content type. + b.Attributes[legacyContentType] = s } return nil } diff --git a/protocol/pubsub/v2/write_pubsub_message_test.go b/protocol/pubsub/v2/write_pubsub_message_test.go new file mode 100644 index 000000000..4c11ec6ea --- /dev/null +++ b/protocol/pubsub/v2/write_pubsub_message_test.go @@ -0,0 +1,254 @@ +/* + Copyright 2026 The CloudEvents Authors + SPDX-License-Identifier: Apache-2.0 +*/ + +package pubsub + +import ( + "context" + "encoding/json" + "testing" + + "cloud.google.com/go/pubsub/v2" + "github.com/stretchr/testify/require" + + "github.com/cloudevents/sdk-go/v2/binding" + "github.com/cloudevents/sdk-go/v2/binding/spec" + bindingtest "github.com/cloudevents/sdk-go/v2/binding/test" + "github.com/cloudevents/sdk-go/v2/binding/transformer" + "github.com/cloudevents/sdk-go/v2/event" + "github.com/cloudevents/sdk-go/v2/test" +) + +func TestWritePubSubMessageBinary(t *testing.T) { + test.EachEvent(t, test.Events(), func(t *testing.T, eventIn event.Event) { + eventIn = test.ConvertEventExtensionsToString(t, eventIn) + ctx := binding.WithForceBinary(context.Background()) + pm := &pubsub.Message{Attributes: make(map[string]string)} + + require.NoError(t, WritePubSubMessage(ctx, binding.ToMessage(&eventIn), pm)) + require.NotContains(t, pm.Attributes, "content-type") + if eventIn.DataContentType() == "" { + require.NotContains(t, pm.Attributes, "ce-datacontenttype") + require.NotContains(t, pm.Attributes, "Content-Type") + } else { + require.Equal(t, eventIn.DataContentType(), pm.Attributes["ce-datacontenttype"]) + require.Equal(t, eventIn.DataContentType(), pm.Attributes["Content-Type"]) + } + require.Equal(t, eventIn.Data(), pm.Data) + + messageOut := NewMessage(pm) + require.Equal(t, binding.EncodingBinary, messageOut.ReadEncoding()) + _, dataContentType := messageOut.GetAttribute(spec.DataContentType) + require.Equal(t, eventIn.DataContentType(), dataContentType) + eventOut, err := binding.ToEvent(ctx, messageOut) + require.NoError(t, err) + test.AssertEventEquals(t, eventIn, *eventOut) + }) +} + +func TestWritePubSubMessageStructured(t *testing.T) { + test.EachEvent(t, test.Events(), func(t *testing.T, eventIn event.Event) { + eventIn = test.ConvertEventExtensionsToString(t, eventIn) + pm := &pubsub.Message{} + ctx := binding.WithForceStructured(context.Background()) + require.NoError(t, WritePubSubMessage(ctx, binding.ToMessage(&eventIn), pm)) + require.Equal(t, event.ApplicationCloudEventsJSON, pm.Attributes["content-type"]) + require.NotContains(t, pm.Attributes, "ce-datacontenttype") + var body map[string]interface{} + require.NoError(t, json.Unmarshal(pm.Data, &body)) + require.NotContains(t, body, "ce-datacontenttype") + if eventIn.DataContentType() == "" { + require.NotContains(t, body, "datacontenttype") + } else { + require.Equal(t, eventIn.DataContentType(), body["datacontenttype"]) + } + msg := NewMessage(pm) + require.Equal(t, binding.EncodingStructured, msg.ReadEncoding()) + eventOut, err := binding.ToEvent(ctx, msg) + require.NoError(t, err) + test.AssertEventEquals(t, eventIn, *eventOut) + }) +} + +func TestWritePubSubMessageChangeEncoding(t *testing.T) { + eventIn := test.ConvertEventExtensionsToString(t, test.FullEvent()) + pm := &pubsub.Message{Attributes: map[string]string{"custom": "preserved"}} + for _, encoding := range []binding.Encoding{binding.EncodingBinary, binding.EncodingStructured, binding.EncodingBinary} { + t.Run(encoding.String(), func(t *testing.T) { + ctx := binding.WithPreferredEventEncoding(context.Background(), encoding) + require.NoError(t, WritePubSubMessage(ctx, binding.ToMessage(&eventIn), pm)) + require.Equal(t, "preserved", pm.Attributes["custom"]) + msg := NewMessage(pm) + require.Equal(t, encoding, msg.ReadEncoding()) + eventOut, err := binding.ToEvent(ctx, msg) + require.NoError(t, err) + test.AssertEventEquals(t, eventIn, *eventOut) + }) + } +} + +func TestWritePubSubMessageReuseDataContentType(t *testing.T) { + tests := []struct { + name string + dataContentType string + }{ + {name: "different data content type", dataContentType: "application/json"}, + {name: "no data content type"}, + } + test.EachEvent(t, test.AllVersions([]event.Event{test.FullEvent()}), func(t *testing.T, previous event.Event) { + previous = test.ConvertEventExtensionsToString(t, previous) + for _, tt := range tests { + for _, encoding := range []binding.Encoding{binding.EncodingBinary, binding.EncodingStructured} { + t.Run(tt.name+"/"+encoding.String(), func(t *testing.T) { + pm := &pubsub.Message{Attributes: map[string]string{"custom": "preserved"}} + require.NoError(t, WritePubSubMessage(binding.WithForceBinary(context.Background()), binding.ToMessage(&previous), pm)) + + next := previous.Clone() + next.SetID(previous.ID() + "-next") + next.SetDataContentType(tt.dataContentType) + ctx := binding.WithPreferredEventEncoding(context.Background(), encoding) + require.NoError(t, WritePubSubMessage(ctx, binding.ToMessage(&next), pm)) + + require.Equal(t, "preserved", pm.Attributes["custom"]) + if encoding == binding.EncodingStructured { + require.Equal(t, event.ApplicationCloudEventsJSON, pm.Attributes["content-type"]) + } else { + require.NotContains(t, pm.Attributes, "content-type") + } + if encoding == binding.EncodingStructured || tt.dataContentType == "" { + require.NotContains(t, pm.Attributes, "ce-datacontenttype") + require.NotContains(t, pm.Attributes, "Content-Type") + } else { + require.Equal(t, tt.dataContentType, pm.Attributes["ce-datacontenttype"]) + require.Equal(t, tt.dataContentType, pm.Attributes["Content-Type"]) + } + + msg := NewMessage(pm) + require.Equal(t, encoding, msg.ReadEncoding()) + eventOut, err := binding.ToEvent(ctx, msg) + require.NoError(t, err) + test.AssertEventEquals(t, next, *eventOut) + }) + } + } + }) +} + +func TestWritePubSubMessageReuseWithMinimalEvent(t *testing.T) { + test.EachEvent(t, test.AllVersions([]event.Event{test.MinEvent()}), func(t *testing.T, next event.Event) { + previous := test.ConvertEventExtensionsToString(t, test.FullEvent()) + previous.Context = spec.VS.Version(next.SpecVersion()).Convert(previous.Context) + for _, encoding := range []binding.Encoding{binding.EncodingBinary, binding.EncodingStructured} { + t.Run(encoding.String(), func(t *testing.T) { + pm := &pubsub.Message{Attributes: map[string]string{"custom": "preserved"}} + require.NoError(t, WritePubSubMessage(binding.WithForceBinary(context.Background()), binding.ToMessage(&previous), pm)) + + ctx := binding.WithPreferredEventEncoding(context.Background(), encoding) + require.NoError(t, WritePubSubMessage(ctx, binding.ToMessage(&next), pm)) + + wantAttributes := map[string]string{"custom": "preserved"} + if encoding == binding.EncodingStructured { + wantAttributes["content-type"] = event.ApplicationCloudEventsJSON + } else { + wantAttributes["ce-specversion"] = next.SpecVersion() + wantAttributes["ce-id"] = next.ID() + wantAttributes["ce-source"] = next.Source() + wantAttributes["ce-type"] = next.Type() + require.Empty(t, pm.Data) + } + require.Equal(t, wantAttributes, pm.Attributes) + + msg := NewMessage(pm) + require.Equal(t, encoding, msg.ReadEncoding()) + eventOut, err := binding.ToEvent(ctx, msg) + require.NoError(t, err) + test.AssertEventEquals(t, next, *eventOut) + }) + } + }) +} + +func TestWritePubSubMessageExistingDataContentType(t *testing.T) { + tests := []struct { + name string + attributes map[string]string + }{ + { + name: "legacy attribute only", + attributes: map[string]string{"Content-Type": "application/protobuf"}, + }, + { + name: "prefixed attribute only", + attributes: map[string]string{"ce-datacontenttype": "application/protobuf"}, + }, + { + name: "legacy and prefixed attributes with matching values", + attributes: map[string]string{ + "Content-Type": "application/protobuf", + "ce-datacontenttype": "application/protobuf", + }, + }, + { + name: "legacy and prefixed attributes with conflicting values", + attributes: map[string]string{ + "Content-Type": "application/protobuf", + "ce-datacontenttype": "text/plain", + }, + }, + } + test.EachEvent(t, test.AllVersions([]event.Event{test.FullEvent()}), func(t *testing.T, eventIn event.Event) { + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + eventIn := test.ConvertEventExtensionsToString(t, eventIn) + ctx := binding.WithForceBinary(context.Background()) + pm := &pubsub.Message{Attributes: map[string]string{"custom": "preserved"}} + for name, value := range tt.attributes { + pm.Attributes[name] = value + } + + require.NoError(t, WritePubSubMessage(ctx, binding.ToMessage(&eventIn), pm)) + require.Equal(t, eventIn.DataContentType(), pm.Attributes["ce-datacontenttype"]) + require.Equal(t, eventIn.DataContentType(), pm.Attributes["Content-Type"]) + require.Equal(t, "preserved", pm.Attributes["custom"]) + require.Equal(t, eventIn.Data(), pm.Data) + + eventOut, err := binding.ToEvent(ctx, NewMessage(pm)) + require.NoError(t, err) + test.AssertEventEquals(t, eventIn, *eventOut) + }) + } + }) +} + +func TestWritePubSubMessageTransformDataContentType(t *testing.T) { + tests := []struct { + name string + value interface{} + }{ + {name: "update", value: "application/json; charset=utf-8"}, + {name: "delete", value: nil}, + } + test.EachEvent(t, test.AllVersions([]event.Event{test.FullEvent()}), func(t *testing.T, eventIn event.Event) { + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + ctx := binding.WithForceBinary(context.Background()) + pm := &pubsub.Message{Attributes: make(map[string]string)} + err := WritePubSubMessage(ctx, bindingtest.MustCreateMockBinaryMessage(eventIn), pm, + transformer.SetAttribute(spec.DataContentType, func(interface{}) (interface{}, error) { + return tt.value, nil + })) + require.NoError(t, err) + if tt.value == nil { + require.NotContains(t, pm.Attributes, "ce-datacontenttype") + require.NotContains(t, pm.Attributes, "Content-Type") + } else { + require.Equal(t, tt.value, pm.Attributes["ce-datacontenttype"]) + require.Equal(t, tt.value, pm.Attributes["Content-Type"]) + } + require.Equal(t, eventIn.Data(), pm.Data) + }) + } + }) +}