Skip to content
Open
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
18 changes: 18 additions & 0 deletions protocol/pubsub/v2/doc.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
75 changes: 59 additions & 16 deletions protocol/pubsub/v2/message.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ package pubsub
import (
"bytes"
"context"
"fmt"
"strings"

"cloud.google.com/go/pubsub/v2"
Expand All @@ -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)
Expand All @@ -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)
}
Expand All @@ -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) {
Expand All @@ -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
Expand All @@ -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]
}
Expand Down
Loading
Loading