From b2f6216f6b3a3f494c0e591259a4e7ce92422d2f Mon Sep 17 00:00:00 2001 From: akme Date: Fri, 17 Apr 2026 11:55:48 +0300 Subject: [PATCH] feat(kafka): add franz-go protocol binding --- docs/index.md | 1 + docs/protocol_implementations.md | 1 + protocol/kafka_franz/v2/doc.go | 7 + protocol/kafka_franz/v2/go.mod | 27 ++ protocol/kafka_franz/v2/go.sum | 50 +++ protocol/kafka_franz/v2/message.go | 183 ++++++++++ protocol/kafka_franz/v2/message_test.go | 132 ++++++++ protocol/kafka_franz/v2/option.go | 75 ++++ protocol/kafka_franz/v2/protocol.go | 320 ++++++++++++++++++ protocol/kafka_franz/v2/protocol_test.go | 315 +++++++++++++++++ .../kafka_franz/v2/write_producer_message.go | 113 +++++++ .../v2/write_producer_message_test.go | 80 +++++ 12 files changed, 1304 insertions(+) create mode 100644 protocol/kafka_franz/v2/doc.go create mode 100644 protocol/kafka_franz/v2/go.mod create mode 100644 protocol/kafka_franz/v2/go.sum create mode 100644 protocol/kafka_franz/v2/message.go create mode 100644 protocol/kafka_franz/v2/message_test.go create mode 100644 protocol/kafka_franz/v2/option.go create mode 100644 protocol/kafka_franz/v2/protocol.go create mode 100644 protocol/kafka_franz/v2/protocol_test.go create mode 100644 protocol/kafka_franz/v2/write_producer_message.go create mode 100644 protocol/kafka_franz/v2/write_producer_message_test.go diff --git a/docs/index.md b/docs/index.md index 44db1b1a0..fbdd157eb 100644 --- a/docs/index.md +++ b/docs/index.md @@ -126,6 +126,7 @@ err := json.Unmarshal(bytes, &event) | [JSON Event Format](event_data_structure.md#marshalunmarshal-event-to-json) | :heavy_check_mark: | :heavy_check_mark: | | [Sarama Kafka Protocol Binding](https://github.com/cloudevents/sdk-go/tree/main/samples/kafka) | :heavy_check_mark: | :heavy_check_mark: | | [Confluent Kafka Protocol Binding](https://github.com/cloudevents/sdk-go/tree/main/samples/kafka_confluent) | :heavy_check_mark: | :heavy_check_mark: | +| [Franz-go Kafka Protocol Binding](https://github.com/cloudevents/sdk-go/tree/main/protocol/kafka_franz) | :heavy_check_mark: | :heavy_check_mark: | | [MQTT Protocol Binding](https://github.com/cloudevents/sdk-go/tree/main/samples/mqtt) | :heavy_check_mark: | :heavy_check_mark: | | [NATS Protocol Binding](https://github.com/cloudevents/sdk-go/tree/main/samples/nats) | :heavy_check_mark: | :heavy_check_mark: | | [STAN Protocol Binding](https://github.com/cloudevents/sdk-go/tree/main/samples/stan) | :heavy_check_mark: | :heavy_check_mark: | diff --git a/docs/protocol_implementations.md b/docs/protocol_implementations.md index 5d154b11b..53bacb530 100644 --- a/docs/protocol_implementations.md +++ b/docs/protocol_implementations.md @@ -26,6 +26,7 @@ the `Write` functions, while the latter is done implementing spec * [Kafka Protocol Binding](https://github.com/cloudevents/sdk-go/tree/main/protocol/kafka_sarama) using [Sarama](https://github.com/Shopify/sarama) * [Kafka Protocol Binding](https://github.com/cloudevents/sdk-go/tree/main/protocol/kafka_confluent) using [confluent-kafka-go](https://github.com/confluentinc/confluent-kafka-go) > It bring some new features compared the above [sarama](https://github.com/Shopify/sarama) binding. Like [pattern subscription](https://github.com/confluentinc/confluent-kafka-go/issues/96), [async message confirmation](https://github.com/cloudevents/sdk-go/issues/846) and other enhancements. +* [Kafka Protocol Binding](https://github.com/cloudevents/sdk-go/tree/main/protocol/kafka_franz) using [franz-go](https://github.com/twmb/franz-go) * [MQTT Protocol Binding](https://github.com/cloudevents/sdk-go/tree/main/protocol/mqtt_paho) using [eclipse/paho.golang](https://github.com/eclipse/paho.golang) * [NATS Protocol Binding](https://github.com/cloudevents/sdk-go/tree/main/protocol/nats) using [nats.go](https://github.com/nats-io/nats.go) * [STAN Protocol Binding](https://github.com/cloudevents/sdk-go/tree/main/protocol/stan) using [stan.go](https://github.com/nats-io/stan.go) diff --git a/protocol/kafka_franz/v2/doc.go b/protocol/kafka_franz/v2/doc.go new file mode 100644 index 000000000..0fd91b52b --- /dev/null +++ b/protocol/kafka_franz/v2/doc.go @@ -0,0 +1,7 @@ +/* + Copyright 2026 The CloudEvents Authors + SPDX-License-Identifier: Apache-2.0 +*/ + +// Package kafka_franz implements a Kafka binding using github.com/twmb/franz-go/pkg/kgo. +package kafka_franz diff --git a/protocol/kafka_franz/v2/go.mod b/protocol/kafka_franz/v2/go.mod new file mode 100644 index 000000000..e3081ef11 --- /dev/null +++ b/protocol/kafka_franz/v2/go.mod @@ -0,0 +1,27 @@ +module github.com/cloudevents/sdk-go/protocol/kafka_franz/v2 + +go 1.25.0 + +replace github.com/cloudevents/sdk-go/v2 => ../../../v2 + +require ( + github.com/cloudevents/sdk-go/v2 v2.16.2 + github.com/stretchr/testify v1.11.1 + github.com/twmb/franz-go v1.20.3 +) + +require ( + github.com/davecgh/go-spew v1.1.1 // indirect + github.com/google/go-cmp v0.7.0 // indirect + github.com/google/uuid v1.6.0 // indirect + github.com/json-iterator/go v1.1.12 // indirect + github.com/klauspost/compress v1.18.0 // indirect + github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect + github.com/modern-go/reflect2 v1.0.2 // indirect + github.com/pierrec/lz4/v4 v4.1.22 // indirect + github.com/pmezard/go-difflib v1.0.0 // indirect + github.com/twmb/franz-go/pkg/kmsg v1.12.0 // indirect + go.uber.org/multierr v1.11.0 // indirect + go.uber.org/zap v1.27.1 // indirect + gopkg.in/yaml.v3 v3.0.1 // indirect +) diff --git a/protocol/kafka_franz/v2/go.sum b/protocol/kafka_franz/v2/go.sum new file mode 100644 index 000000000..ac08242fe --- /dev/null +++ b/protocol/kafka_franz/v2/go.sum @@ -0,0 +1,50 @@ +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM= +github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo= +github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo= +github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= +github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd h1:TRLaZ9cD/w8PVh93nsPXa1VrQ6jlwL5oN8l14QlcNfg= +github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= +github.com/modern-go/reflect2 v1.0.2 h1:xBagoLtFs94CBntxluKeaWgTMpvLxC4ur3nMaC9Gz0M= +github.com/modern-go/reflect2 v1.0.2/go.mod h1:yWuevngMOJpCy52FWWMvUC8ws7m/LJsjYzDa0/r8luk= +github.com/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e h1:fD57ERR4JtEqsWbfPhv4DMiApHyliiK5xCTNVSPiaAs= +github.com/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e/go.mod h1:zD1mROLANZcx1PVRCS0qkT7pwLkGfwJo4zjcN/Tysno= +github.com/pierrec/lz4/v4 v4.1.22 h1:cKFw6uJDK+/gfw5BcDL0JL5aBsAFdsIT18eRtLj7VIU= +github.com/pierrec/lz4/v4 v4.1.22/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +github.com/twmb/franz-go v1.20.3 h1:gjwZwZmmvo/t7mxyj6frxDORVxsqrycXPnDrpkXldfY= +github.com/twmb/franz-go v1.20.3/go.mod h1:YCnepDd4gl6vdzG03I5Wa57RnCTIC6DVEyMpDX/J8UA= +github.com/twmb/franz-go/pkg/kmsg v1.12.0 h1:CbatD7ers1KzDNgJqPbKOq0Bz/WLBdsTH75wgzeVaPc= +github.com/twmb/franz-go/pkg/kmsg v1.12.0/go.mod h1:+DPt4NC8RmI6hqb8G09+3giKObE6uD2Eya6CfqBpeJY= +github.com/valyala/bytebufferpool v1.0.0 h1:GqA5TC/0021Y/b9FG4Oi9Mr3q7XYx6KllzawFIhcdPw= +github.com/valyala/bytebufferpool v1.0.0/go.mod h1:6bBcMArwyJ5K/AmCkWv1jt77kVWyCJ6HpOuEn7z0Csc= +go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= +go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= +go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0= +go.uber.org/multierr v1.11.0/go.mod h1:20+QtiLqy0Nd6FdQB9TLXag12DsQkrbs3htMFfDN80Y= +go.uber.org/zap v1.27.1 h1:08RqriUEv8+ArZRYSTXy1LeBScaMpVSTBhCeaZYfMYc= +go.uber.org/zap v1.27.1/go.mod h1:GB2qFLM7cTU87MWRP2mPIjqfIDnGu+VIO4V/SdhGo2E= +golang.org/x/crypto v0.43.0 h1:dduJYIi3A3KOfdGOHX8AVZ/jGiyPa3IbBozJ5kNuE04= +golang.org/x/crypto v0.43.0/go.mod h1:BFbav4mRNlXJL4wNeejLpWxB7wMbc79PdRGhWKncxR0= +golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U= +golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20200227125254-8fa46927fb4f h1:BLraFXnmrev5lT+xlilqcH8XK9/i0At2xKjWk4p6zsU= +gopkg.in/check.v1 v1.0.0-20200227125254-8fa46927fb4f/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/protocol/kafka_franz/v2/message.go b/protocol/kafka_franz/v2/message.go new file mode 100644 index 000000000..2ea9c04f7 --- /dev/null +++ b/protocol/kafka_franz/v2/message.go @@ -0,0 +1,183 @@ +/* + Copyright 2026 The CloudEvents Authors + SPDX-License-Identifier: Apache-2.0 +*/ + +package kafka_franz + +import ( + "bytes" + "context" + "errors" + "strconv" + "strings" + "sync" + + "github.com/twmb/franz-go/pkg/kgo" + + "github.com/cloudevents/sdk-go/v2/binding" + "github.com/cloudevents/sdk-go/v2/binding/format" + "github.com/cloudevents/sdk-go/v2/binding/spec" +) + +const ( + prefix = "ce_" + contentTypeHeader = "content-type" +) + +const ( + KafkaOffsetKey = "kafkaoffset" + KafkaPartitionKey = "kafkapartition" + KafkaTopicKey = "kafkatopic" + KafkaMessageKey = "kafkamessagekey" +) + +var specs = spec.WithPrefix(prefix) + +// Message represents a Kafka message. +// This message can be read several times safely. +type Message struct { + record *kgo.Record + properties map[string][]byte + format format.Format + version spec.Version + contentType string +} + +var ( + _ binding.Message = (*Message)(nil) + _ binding.MessageMetadataReader = (*Message)(nil) +) + +// NewMessage returns a binding.Message that holds the provided kgo.Record. +// The returned binding.Message can be read several times safely. +func NewMessage(record *kgo.Record) *Message { + if record == nil { + panic("the kgo.Record must not be nil") + } + + var contentType, contentVersion string + properties := make(map[string][]byte, len(record.Headers)+4) + for _, header := range record.Headers { + key := strings.ToLower(header.Key) + if key == contentTypeHeader { + contentType = string(header.Value) + } + if key == specs.PrefixedSpecVersionName() { + contentVersion = string(header.Value) + } + properties[key] = append([]byte(nil), header.Value...) + } + + properties[prefix+KafkaOffsetKey] = []byte(strconv.FormatInt(record.Offset, 10)) + properties[prefix+KafkaPartitionKey] = []byte(strconv.FormatInt(int64(record.Partition), 10)) + properties[prefix+KafkaTopicKey] = []byte(record.Topic) + if len(record.Key) > 0 { + properties[prefix+KafkaMessageKey] = append([]byte(nil), record.Key...) + } + + message := &Message{ + record: record, + properties: properties, + contentType: contentType, + } + if ft := format.Lookup(contentType); ft != nil { + message.format = ft + } else if v := specs.Version(contentVersion); v != nil { + message.version = v + } + + return message +} + +func (m *Message) ReadEncoding() binding.Encoding { + if m.version != nil { + return binding.EncodingBinary + } + if m.format != nil { + return binding.EncodingStructured + } + return binding.EncodingUnknown +} + +func (m *Message) ReadStructured(ctx context.Context, encoder binding.StructuredWriter) error { + if m.format == nil { + return binding.ErrNotStructured + } + return encoder.SetStructuredEvent(ctx, m.format, bytes.NewReader(m.record.Value)) +} + +func (m *Message) ReadBinary(ctx context.Context, encoder binding.BinaryWriter) error { + if m.version == nil { + return binding.ErrNotBinary + } + + var err error + for key, value := range m.properties { + switch { + case strings.HasPrefix(key, prefix): + attr := m.version.Attribute(key) + if attr != nil { + err = encoder.SetAttribute(attr, string(value)) + } else { + err = encoder.SetExtension(strings.TrimPrefix(key, prefix), string(value)) + } + case key == contentTypeHeader: + err = encoder.SetAttribute(m.version.AttributeFromKind(spec.DataContentType), string(value)) + } + if err != nil { + return err + } + } + + if m.record.Value != nil { + err = encoder.SetData(bytes.NewBuffer(m.record.Value)) + } + return err +} + +func (m *Message) Finish(error) error { + return nil +} + +func (m *Message) GetAttribute(k spec.Kind) (spec.Attribute, interface{}) { + if m.version == nil { + return nil, nil + } + attr := m.version.AttributeFromKind(k) + if attr == nil { + return nil, nil + } + return attr, string(m.properties[attr.PrefixedName()]) +} + +func (m *Message) GetExtension(name string) interface{} { + value, ok := m.properties[prefix+name] + if !ok { + return nil + } + return string(value) +} + +type receivedMessage struct { + *Message + finish func(error) error + finishOnce sync.Once + finishErr error +} + +var _ binding.MessageWrapper = (*receivedMessage)(nil) + +func (m *receivedMessage) Finish(err error) error { + m.finishOnce.Do(func() { + m.finishErr = m.Message.Finish(err) + if m.finish != nil { + m.finishErr = errors.Join(m.finishErr, m.finish(err)) + } + }) + return m.finishErr +} + +func (m *receivedMessage) GetWrappedMessage() binding.Message { + return m.Message +} diff --git a/protocol/kafka_franz/v2/message_test.go b/protocol/kafka_franz/v2/message_test.go new file mode 100644 index 000000000..3fcff3b40 --- /dev/null +++ b/protocol/kafka_franz/v2/message_test.go @@ -0,0 +1,132 @@ +/* + Copyright 2026 The CloudEvents Authors + SPDX-License-Identifier: Apache-2.0 +*/ + +package kafka_franz + +import ( + "context" + "testing" + + "github.com/stretchr/testify/require" + "github.com/twmb/franz-go/pkg/kgo" + + cloudevents "github.com/cloudevents/sdk-go/v2" + "github.com/cloudevents/sdk-go/v2/binding" + "github.com/cloudevents/sdk-go/v2/binding/format" + "github.com/cloudevents/sdk-go/v2/test" +) + +var ( + ctx = context.Background() + testEvent = test.FullEvent() + + structuredConsumerRecord = &kgo.Record{ + Topic: "test-topic", + Partition: 0, + Offset: 10, + Value: func() []byte { + b, _ := format.JSON.Marshal(&testEvent) + return b + }(), + Headers: []kgo.RecordHeader{{ + Key: contentTypeHeader, + Value: []byte(cloudevents.ApplicationCloudEventsJSON), + }}, + } + + binaryConsumerRecord = &kgo.Record{ + Topic: "test-topic", + Partition: 0, + Offset: 10, + Value: []byte("hello world!"), + Headers: mapToRecordHeaders(map[string]string{ + "ce_type": testEvent.Type(), + "ce_source": testEvent.Source(), + "ce_id": testEvent.ID(), + "ce_time": test.Timestamp.String(), + "ce_specversion": "1.0", + "ce_dataschema": test.Schema.String(), + "ce_datacontenttype": "text/json", + "ce_subject": "receiverTopic", + "exta": "someext", + }), + } +) + +func TestNewMessage(t *testing.T) { + tests := []struct { + name string + record *kgo.Record + expectedEncoding binding.Encoding + }{ + { + name: "structured encoding", + record: structuredConsumerRecord, + expectedEncoding: binding.EncodingStructured, + }, + { + name: "binary encoding", + record: binaryConsumerRecord, + expectedEncoding: binding.EncodingBinary, + }, + { + name: "unknown encoding", + record: &kgo.Record{ + Topic: "test-topic", + Partition: 0, + Offset: 10, + Value: []byte("{}"), + Headers: []kgo.RecordHeader{{ + Key: contentTypeHeader, + Value: []byte("application/json"), + }}, + }, + expectedEncoding: binding.EncodingUnknown, + }, + { + name: "binary encoding with empty value", + record: &kgo.Record{ + Topic: "test-topic", + Partition: 0, + Offset: 10, + Headers: mapToRecordHeaders(map[string]string{ + "ce_type": testEvent.Type(), + "ce_source": testEvent.Source(), + "ce_id": testEvent.ID(), + "ce_time": test.Timestamp.String(), + "ce_specversion": "1.0", + "ce_dataschema": test.Schema.String(), + "ce_datacontenttype": "text/json", + "ce_subject": "receiverTopic", + }), + }, + expectedEncoding: binding.EncodingBinary, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + msg := NewMessage(tt.record) + require.Equal(t, tt.expectedEncoding, msg.ReadEncoding()) + + var err error + switch tt.expectedEncoding { + case binding.EncodingStructured: + err = msg.ReadStructured(ctx, (*kafkaRecordWriter)(tt.record)) + case binding.EncodingBinary: + err = msg.ReadBinary(ctx, (*kafkaRecordWriter)(tt.record)) + } + require.NoError(t, err) + }) + } +} + +func mapToRecordHeaders(m map[string]string) []kgo.RecordHeader { + res := make([]kgo.RecordHeader, 0, len(m)) + for k, v := range m { + res = append(res, kgo.RecordHeader{Key: k, Value: []byte(v)}) + } + return res +} diff --git a/protocol/kafka_franz/v2/option.go b/protocol/kafka_franz/v2/option.go new file mode 100644 index 000000000..00fba255f --- /dev/null +++ b/protocol/kafka_franz/v2/option.go @@ -0,0 +1,75 @@ +/* + Copyright 2026 The CloudEvents Authors + SPDX-License-Identifier: Apache-2.0 +*/ + +package kafka_franz + +import ( + "context" + "errors" + + "github.com/twmb/franz-go/pkg/kgo" +) + +// Option is the function signature required to be considered a kafka_franz.Option. +type Option func(*Protocol) error + +// WithClient sets a kgo.Client instance to initialize the protocol directly. +func WithClient(client *kgo.Client) Option { + return func(p *Protocol) error { + if client == nil { + return errors.New("the kgo.Client option must not be nil") + } + p.client = newKgoClient(client) + return nil + } +} + +// WithClientOptions sets the options used to create a kgo.Client. +// When the protocol creates the client, it also enables manual offset control +// with kgo.DisableAutoCommit and kgo.BlockRebalanceOnPoll so ACKs drive commits. +func WithClientOptions(opts ...kgo.Opt) Option { + return func(p *Protocol) error { + if len(opts) == 0 { + return errors.New("the kgo client options must not be empty") + } + p.clientOptions = append(p.clientOptions, opts...) + return nil + } +} + +// WithSenderTopic sets the default topic used for produces when the context does not override it. +func WithSenderTopic(topic string) Option { + return func(p *Protocol) error { + if topic == "" { + return errors.New("the producer topic option must not be empty") + } + p.producerDefaultTopic = topic + return nil + } +} + +// Opaque key type used to store the Kafka message key. +type messageKeyType struct{} + +var keyForMessageKey = messageKeyType{} + +// WithMessageKey returns a new context with the given Kafka message key. +func WithMessageKey(ctx context.Context, messageKey []byte) context.Context { + keyCopy := append([]byte(nil), messageKey...) + return context.WithValue(ctx, keyForMessageKey, keyCopy) +} + +// MessageKeyFrom looks in the given context and returns the message key if found, otherwise nil. +func MessageKeyFrom(ctx context.Context) []byte { + c := ctx.Value(keyForMessageKey) + if c == nil { + return nil + } + s, ok := c.([]byte) + if !ok { + return nil + } + return append([]byte(nil), s...) +} diff --git a/protocol/kafka_franz/v2/protocol.go b/protocol/kafka_franz/v2/protocol.go new file mode 100644 index 000000000..e25ced16e --- /dev/null +++ b/protocol/kafka_franz/v2/protocol.go @@ -0,0 +1,320 @@ +/* + Copyright 2026 The CloudEvents Authors + SPDX-License-Identifier: Apache-2.0 +*/ + +package kafka_franz + +import ( + "context" + "errors" + "fmt" + "io" + "sync" + + "github.com/twmb/franz-go/pkg/kgo" + + "github.com/cloudevents/sdk-go/v2/binding" + cecontext "github.com/cloudevents/sdk-go/v2/context" + "github.com/cloudevents/sdk-go/v2/protocol" +) + +var ( + _ protocol.Sender = (*Protocol)(nil) + _ protocol.Opener = (*Protocol)(nil) + _ protocol.Receiver = (*Protocol)(nil) + _ protocol.Closer = (*Protocol)(nil) +) + +type produceResult interface { + FirstErr() error +} + +type fetchResult interface { + IsClientClosed() bool + Errors() []kgo.FetchError + Records() []*kgo.Record +} + +type kafkaClient interface { + ProduceSync(context.Context, ...*kgo.Record) produceResult + PollFetches(context.Context) fetchResult + CommitRecords(context.Context, ...*kgo.Record) error + AllowRebalance() + Close() + CloseAllowingRebalance() +} + +type kgoClient struct { + client *kgo.Client +} + +type kgoProduceResult struct { + result kgo.ProduceResults +} + +type kgoFetchResult struct { + result kgo.Fetches +} + +func newKgoClient(client *kgo.Client) kafkaClient { + return &kgoClient{client: client} +} + +func (c *kgoClient) ProduceSync(ctx context.Context, records ...*kgo.Record) produceResult { + return kgoProduceResult{result: c.client.ProduceSync(ctx, records...)} +} + +func (c *kgoClient) PollFetches(ctx context.Context) fetchResult { + return kgoFetchResult{result: c.client.PollFetches(ctx)} +} + +func (c *kgoClient) CommitRecords(ctx context.Context, records ...*kgo.Record) error { + return c.client.CommitRecords(ctx, records...) +} + +func (c *kgoClient) AllowRebalance() { + c.client.AllowRebalance() +} + +func (c *kgoClient) Close() { + c.client.Close() +} + +func (c *kgoClient) CloseAllowingRebalance() { + c.client.CloseAllowingRebalance() +} + +func (r kgoProduceResult) FirstErr() error { + return r.result.FirstErr() +} + +func (r kgoFetchResult) IsClientClosed() bool { + return r.result.IsClientClosed() +} + +func (r kgoFetchResult) Errors() []kgo.FetchError { + return r.result.Errors() +} + +func (r kgoFetchResult) Records() []*kgo.Record { + return r.result.Records() +} + +// Protocol implements a CloudEvents Kafka transport backed by franz-go. +type Protocol struct { + client kafkaClient + clientOptions []kgo.Opt + ownClient bool + producerDefaultTopic string + + consumerIncoming chan binding.Message + + receiverMux sync.Mutex + receiverCtx context.Context + receiverCancel context.CancelFunc + receiverOpened bool +} + +// New creates a new franz-go Kafka transport. +func New(opts ...Option) (*Protocol, error) { + p := &Protocol{ + consumerIncoming: make(chan binding.Message), + } + if err := p.applyOptions(opts...); err != nil { + return nil, err + } + + if p.client != nil && len(p.clientOptions) > 0 { + return nil, errors.New("the client and kgo client options must not be set together") + } + + if p.client == nil { + if len(p.clientOptions) == 0 { + return nil, errors.New("at least one of the following to initialize the protocol must be set: client or kgo client options") + } + clientOptions := append([]kgo.Opt{}, p.clientOptions...) + clientOptions = append(clientOptions, kgo.DisableAutoCommit(), kgo.BlockRebalanceOnPoll()) + + client, err := kgo.NewClient(clientOptions...) + if err != nil { + return nil, err + } + p.client = newKgoClient(client) + p.ownClient = true + } + + return p, nil +} + +func (p *Protocol) applyOptions(opts ...Option) error { + for _, fn := range opts { + if err := fn(p); err != nil { + return err + } + } + return nil +} + +// Send transmits a CloudEvent using franz-go's synchronous produce API. +func (p *Protocol) Send(ctx context.Context, in binding.Message, transformers ...binding.Transformer) (err error) { + if p.client == nil { + return errors.New("producer client must be set") + } + defer in.Finish(err) + + topic := cecontext.TopicFrom(ctx) + if topic == "" { + topic = p.producerDefaultTopic + } + if topic == "" { + return errors.New("the producer topic must be set either by option or context") + } + + record := &kgo.Record{Topic: topic} + if key := MessageKeyFrom(ctx); key != nil { + record.Key = key + } + + if err = WriteProducerMessage(ctx, in, record, transformers...); err != nil { + return fmt.Errorf("create producer record: %w", err) + } + + if err = p.client.ProduceSync(ctx, record).FirstErr(); err != nil { + return fmt.Errorf("produce record: %w", err) + } + return nil +} + +// OpenInbound starts the receive loop. This call blocks until the receiver stops. +func (p *Protocol) OpenInbound(ctx context.Context) error { + if p.client == nil { + return errors.New("consumer client must be set") + } + + p.receiverMux.Lock() + if p.receiverOpened { + p.receiverMux.Unlock() + return errors.New("receiver already open") + } + p.receiverOpened = true + p.receiverCtx, p.receiverCancel = context.WithCancel(ctx) + receiveCtx := p.receiverCtx + p.receiverMux.Unlock() + + defer func() { + p.receiverMux.Lock() + if p.receiverCancel != nil { + p.receiverCancel() + p.receiverCancel = nil + } + p.receiverCtx = nil + p.receiverMux.Unlock() + close(p.consumerIncoming) + }() + + logger := cecontext.LoggerFrom(ctx) + + for { + fetches := p.client.PollFetches(receiveCtx) + if fetches.IsClientClosed() { + return nil + } + + records := fetches.Records() + if len(records) > 0 { + batch := receiveBatch{ctx: receiveCtx, client: p.client} + for _, record := range records { + msg := batch.wrap(record) + select { + case p.consumerIncoming <- msg: + case <-receiveCtx.Done(): + batch.drop() + p.client.AllowRebalance() + return receiveCtx.Err() + } + } + batch.wait() + p.client.AllowRebalance() + } + + if err := joinFetchErrors(fetches.Errors()); err != nil { + logger.Warnw("franz-go fetch error", "error", err) + return err + } + + if err := receiveCtx.Err(); err != nil { + return err + } + } +} + +// Receive implements protocol.Receiver. +func (p *Protocol) Receive(ctx context.Context) (binding.Message, error) { + select { + case msg, ok := <-p.consumerIncoming: + if !ok { + return nil, io.EOF + } + return msg, nil + case <-ctx.Done(): + return nil, io.EOF + } +} + +// Close cleans up resources after use. +func (p *Protocol) Close(context.Context) error { + p.receiverMux.Lock() + cancel := p.receiverCancel + p.receiverCancel = nil + p.receiverCtx = nil + p.receiverMux.Unlock() + + if cancel != nil { + cancel() + } + if p.ownClient && p.client != nil { + p.client.CloseAllowingRebalance() + } + return nil +} + +type receiveBatch struct { + ctx context.Context + client kafkaClient + wg sync.WaitGroup +} + +func (b *receiveBatch) wrap(record *kgo.Record) *receivedMessage { + b.wg.Add(1) + return &receivedMessage{ + Message: NewMessage(record), + finish: func(err error) error { + defer b.wg.Done() + if !protocol.IsACK(err) || b.ctx.Err() != nil { + return nil + } + return b.client.CommitRecords(b.ctx, record) + }, + } +} + +func (b *receiveBatch) drop() { + b.wg.Done() +} + +func (b *receiveBatch) wait() { + b.wg.Wait() +} + +func joinFetchErrors(fetchErrors []kgo.FetchError) error { + if len(fetchErrors) == 0 { + return nil + } + + errs := make([]error, 0, len(fetchErrors)) + for _, fetchErr := range fetchErrors { + errs = append(errs, fmt.Errorf("fetch error for topic %s partition %d: %w", fetchErr.Topic, fetchErr.Partition, fetchErr.Err)) + } + return errors.Join(errs...) +} diff --git a/protocol/kafka_franz/v2/protocol_test.go b/protocol/kafka_franz/v2/protocol_test.go new file mode 100644 index 000000000..c203c5d92 --- /dev/null +++ b/protocol/kafka_franz/v2/protocol_test.go @@ -0,0 +1,315 @@ +/* + Copyright 2026 The CloudEvents Authors + SPDX-License-Identifier: Apache-2.0 +*/ + +package kafka_franz + +import ( + "context" + "errors" + "testing" + "time" + + "github.com/stretchr/testify/require" + "github.com/twmb/franz-go/pkg/kgo" + + "github.com/cloudevents/sdk-go/v2/binding" + cecontext "github.com/cloudevents/sdk-go/v2/context" + "github.com/cloudevents/sdk-go/v2/protocol" + "github.com/cloudevents/sdk-go/v2/test" +) + +type fakeProduceResult struct { + err error +} + +func (r fakeProduceResult) FirstErr() error { + return r.err +} + +type fakeFetchResult struct { + clientClosed bool + errors []kgo.FetchError + records []*kgo.Record +} + +func (r fakeFetchResult) IsClientClosed() bool { + return r.clientClosed +} + +func (r fakeFetchResult) Errors() []kgo.FetchError { + return r.errors +} + +func (r fakeFetchResult) Records() []*kgo.Record { + return r.records +} + +type fakeClient struct { + produceErr error + commitErr error + + producedRecords []*kgo.Record + committedRecords []*kgo.Record + fetches []fetchResult + + allowRebalanceCalls int + closeCalls int + closeAllowCalls int +} + +func (c *fakeClient) ProduceSync(context.Context, ...*kgo.Record) produceResult { + panic("not implemented") +} + +func (c *fakeClient) PollFetches(context.Context) fetchResult { + if len(c.fetches) == 0 { + return fakeFetchResult{clientClosed: true} + } + result := c.fetches[0] + c.fetches = c.fetches[1:] + return result +} + +func (c *fakeClient) CommitRecords(_ context.Context, records ...*kgo.Record) error { + for _, record := range records { + c.committedRecords = append(c.committedRecords, cloneRecord(record)) + } + return c.commitErr +} + +func (c *fakeClient) AllowRebalance() { + c.allowRebalanceCalls++ +} + +func (c *fakeClient) Close() { + c.closeCalls++ +} + +func (c *fakeClient) CloseAllowingRebalance() { + c.closeAllowCalls++ +} + +type fakeProtocolClient struct { + *fakeClient +} + +func (c *fakeProtocolClient) ProduceSync(_ context.Context, records ...*kgo.Record) produceResult { + for _, record := range records { + c.producedRecords = append(c.producedRecords, cloneRecord(record)) + } + return fakeProduceResult{err: c.produceErr} +} + +func TestNewProtocol(t *testing.T) { + t.Run("requires client configuration", func(t *testing.T) { + _, err := New() + require.EqualError(t, err, "at least one of the following to initialize the protocol must be set: client or kgo client options") + }) + + t.Run("client and client options are mutually exclusive", func(t *testing.T) { + _, err := New( + WithClient(&kgo.Client{}), + WithClientOptions(kgo.SeedBrokers("127.0.0.1:9092")), + ) + require.EqualError(t, err, "the client and kgo client options must not be set together") + }) +} + +func TestSendUsesDefaultTopicContextTopicAndMessageKey(t *testing.T) { + client := &fakeProtocolClient{fakeClient: &fakeClient{}} + p := &Protocol{ + client: client, + producerDefaultTopic: "default-topic", + consumerIncoming: make(chan binding.Message), + } + + event := test.FullEvent() + msg := (*binding.EventMessage)(&event) + + ctx := WithMessageKey(cecontext.WithTopic(context.Background(), "override-topic"), []byte("record-key")) + err := p.Send(ctx, msg) + require.NoError(t, err) + + require.Len(t, client.producedRecords, 1) + record := client.producedRecords[0] + require.Equal(t, "override-topic", record.Topic) + require.Equal(t, []byte("record-key"), record.Key) + require.NotEmpty(t, record.Headers) +} + +func TestSendRequiresTopic(t *testing.T) { + client := &fakeProtocolClient{fakeClient: &fakeClient{}} + p := &Protocol{ + client: client, + consumerIncoming: make(chan binding.Message), + } + + event := test.FullEvent() + msg := (*binding.EventMessage)(&event) + + err := p.Send(context.Background(), msg) + require.EqualError(t, err, "the producer topic must be set either by option or context") + require.Empty(t, client.producedRecords) +} + +func TestOpenInboundACKCommitsRecord(t *testing.T) { + record := binaryRecord("input-topic", 2, 42) + client := &fakeProtocolClient{fakeClient: &fakeClient{ + fetches: []fetchResult{ + fakeFetchResult{records: []*kgo.Record{record}}, + fakeFetchResult{clientClosed: true}, + }, + }} + p := &Protocol{ + client: client, + consumerIncoming: make(chan binding.Message), + } + + errCh := make(chan error, 1) + go func() { + errCh <- p.OpenInbound(context.Background()) + }() + + msg, err := p.Receive(context.Background()) + require.NoError(t, err) + require.NoError(t, msg.Finish(protocol.ResultACK)) + require.NoError(t, waitForError(t, errCh)) + + require.Len(t, client.committedRecords, 1) + require.Equal(t, record.Topic, client.committedRecords[0].Topic) + require.Equal(t, record.Offset, client.committedRecords[0].Offset) + require.Equal(t, 1, client.allowRebalanceCalls) +} + +func TestOpenInboundNACKSkipsCommit(t *testing.T) { + record := binaryRecord("input-topic", 0, 5) + client := &fakeProtocolClient{fakeClient: &fakeClient{ + fetches: []fetchResult{ + fakeFetchResult{records: []*kgo.Record{record}}, + fakeFetchResult{clientClosed: true}, + }, + }} + p := &Protocol{ + client: client, + consumerIncoming: make(chan binding.Message), + } + + errCh := make(chan error, 1) + go func() { + errCh <- p.OpenInbound(context.Background()) + }() + + msg, err := p.Receive(context.Background()) + require.NoError(t, err) + require.NoError(t, msg.Finish(protocol.ResultNACK)) + require.NoError(t, waitForError(t, errCh)) + + require.Empty(t, client.committedRecords) + require.Equal(t, 1, client.allowRebalanceCalls) +} + +func TestOpenInboundFinishReturnsCommitError(t *testing.T) { + commitErr := errors.New("commit failed") + client := &fakeProtocolClient{fakeClient: &fakeClient{ + commitErr: commitErr, + fetches: []fetchResult{ + fakeFetchResult{records: []*kgo.Record{binaryRecord("input-topic", 0, 7)}}, + fakeFetchResult{clientClosed: true}, + }, + }} + p := &Protocol{ + client: client, + consumerIncoming: make(chan binding.Message), + } + + errCh := make(chan error, 1) + go func() { + errCh <- p.OpenInbound(context.Background()) + }() + + msg, err := p.Receive(context.Background()) + require.NoError(t, err) + require.EqualError(t, msg.Finish(protocol.ResultACK), "commit failed") + require.NoError(t, waitForError(t, errCh)) +} + +func TestOpenInboundReturnsFetchErrors(t *testing.T) { + client := &fakeProtocolClient{fakeClient: &fakeClient{ + fetches: []fetchResult{ + fakeFetchResult{ + errors: []kgo.FetchError{{ + Topic: "input-topic", + Partition: 3, + Err: errors.New("permission denied"), + }}, + }, + }, + }} + p := &Protocol{ + client: client, + consumerIncoming: make(chan binding.Message), + } + + err := p.OpenInbound(context.Background()) + require.EqualError(t, err, "fetch error for topic input-topic partition 3: permission denied") +} + +func TestCloseOwnClient(t *testing.T) { + client := &fakeProtocolClient{fakeClient: &fakeClient{}} + p := &Protocol{ + client: client, + ownClient: true, + consumerIncoming: make(chan binding.Message), + } + + require.NoError(t, p.Close(context.Background())) + require.Equal(t, 1, client.closeAllowCalls) + require.Zero(t, client.closeCalls) +} + +func waitForError(t *testing.T, errCh <-chan error) error { + t.Helper() + select { + case err := <-errCh: + return err + case <-time.After(2 * time.Second): + t.Fatal("timed out waiting for OpenInbound to finish") + return nil + } +} + +func binaryRecord(topic string, partition int32, offset int64) *kgo.Record { + return &kgo.Record{ + Topic: topic, + Partition: partition, + Offset: offset, + Value: []byte("hello"), + Headers: []kgo.RecordHeader{ + {Key: "ce_specversion", Value: []byte("1.0")}, + {Key: "ce_type", Value: []byte("example.type")}, + {Key: "ce_source", Value: []byte("example/source")}, + {Key: "ce_id", Value: []byte("example-id")}, + {Key: "ce_datacontenttype", Value: []byte("text/plain")}, + }, + } +} + +func cloneRecord(record *kgo.Record) *kgo.Record { + if record == nil { + return nil + } + + cloned := *record + cloned.Key = append([]byte(nil), record.Key...) + cloned.Value = append([]byte(nil), record.Value...) + cloned.Headers = make([]kgo.RecordHeader, len(record.Headers)) + for i, header := range record.Headers { + cloned.Headers[i] = kgo.RecordHeader{ + Key: header.Key, + Value: append([]byte(nil), header.Value...), + } + } + return &cloned +} diff --git a/protocol/kafka_franz/v2/write_producer_message.go b/protocol/kafka_franz/v2/write_producer_message.go new file mode 100644 index 000000000..1d951f599 --- /dev/null +++ b/protocol/kafka_franz/v2/write_producer_message.go @@ -0,0 +1,113 @@ +/* + Copyright 2026 The CloudEvents Authors + SPDX-License-Identifier: Apache-2.0 +*/ + +package kafka_franz + +import ( + "bytes" + "context" + "io" + + "github.com/twmb/franz-go/pkg/kgo" + + "github.com/cloudevents/sdk-go/v2/binding" + "github.com/cloudevents/sdk-go/v2/binding/format" + "github.com/cloudevents/sdk-go/v2/binding/spec" + "github.com/cloudevents/sdk-go/v2/types" +) + +type kafkaRecordWriter kgo.Record + +var ( + _ binding.StructuredWriter = (*kafkaRecordWriter)(nil) + _ binding.BinaryWriter = (*kafkaRecordWriter)(nil) +) + +// WriteProducerMessage fills the provided record with the message in. +func WriteProducerMessage(ctx context.Context, in binding.Message, record *kgo.Record, transformers ...binding.Transformer) error { + writer := (*kafkaRecordWriter)(record) + _, err := binding.Write(ctx, in, writer, writer, transformers...) + return err +} + +func (w *kafkaRecordWriter) SetStructuredEvent(ctx context.Context, f format.Format, event io.Reader) error { + w.Headers = []kgo.RecordHeader{{ + Key: contentTypeHeader, + Value: []byte(f.MediaType()), + }} + + var buf bytes.Buffer + if _, err := io.Copy(&buf, event); err != nil { + return err + } + + w.Value = buf.Bytes() + return nil +} + +func (w *kafkaRecordWriter) Start(context.Context) error { + w.Headers = []kgo.RecordHeader{} + return nil +} + +func (w *kafkaRecordWriter) End(context.Context) error { + return nil +} + +func (w *kafkaRecordWriter) SetData(reader io.Reader) error { + buf, ok := reader.(*bytes.Buffer) + if !ok { + buf = new(bytes.Buffer) + if _, err := io.Copy(buf, reader); err != nil { + return err + } + } + w.Value = buf.Bytes() + return nil +} + +func (w *kafkaRecordWriter) SetAttribute(attribute spec.Attribute, value interface{}) error { + if attribute.Kind() == spec.DataContentType { + if value == nil { + w.removeHeader(contentTypeHeader) + return nil + } + return w.addHeader(contentTypeHeader, value) + } + + key := prefix + attribute.Name() + if value == nil { + w.removeHeader(key) + return nil + } + return w.addHeader(key, value) +} + +func (w *kafkaRecordWriter) SetExtension(name string, value interface{}) error { + key := prefix + name + if value == nil { + w.removeHeader(key) + return nil + } + return w.addHeader(key, value) +} + +func (w *kafkaRecordWriter) removeHeader(key string) { + for i, header := range w.Headers { + if header.Key == key { + w.Headers = append(w.Headers[:i], w.Headers[i+1:]...) + return + } + } +} + +func (w *kafkaRecordWriter) addHeader(key string, value interface{}) error { + s, err := types.Format(value) + if err != nil { + return err + } + w.Headers = append(w.Headers, kgo.RecordHeader{Key: key, Value: []byte(s)}) + return nil +} diff --git a/protocol/kafka_franz/v2/write_producer_message_test.go b/protocol/kafka_franz/v2/write_producer_message_test.go new file mode 100644 index 000000000..a4392dc1b --- /dev/null +++ b/protocol/kafka_franz/v2/write_producer_message_test.go @@ -0,0 +1,80 @@ +/* + Copyright 2026 The CloudEvents Authors + SPDX-License-Identifier: Apache-2.0 +*/ + +package kafka_franz + +import ( + "context" + "strconv" + "testing" + + "github.com/stretchr/testify/require" + "github.com/twmb/franz-go/pkg/kgo" + + "github.com/cloudevents/sdk-go/v2/binding" + . "github.com/cloudevents/sdk-go/v2/binding/test" + "github.com/cloudevents/sdk-go/v2/event" + . "github.com/cloudevents/sdk-go/v2/test" +) + +func TestWriteProducerMessage(t *testing.T) { + tests := []struct { + name string + context context.Context + messageFactory func(e event.Event) binding.Message + expectedEncoding binding.Encoding + }{ + { + name: "structured to structured", + context: ctx, + messageFactory: func(e event.Event) binding.Message { + return MustCreateMockStructuredMessage(t, e) + }, + expectedEncoding: binding.EncodingStructured, + }, + { + name: "binary to binary", + context: ctx, + messageFactory: MustCreateMockBinaryMessage, + expectedEncoding: binding.EncodingBinary, + }, + } + + EachEvent(t, Events(), func(t *testing.T, e event.Event) { + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + record := &kgo.Record{ + Topic: "test-topic", + Partition: 0, + Offset: 10, + } + + eventIn := ConvertEventExtensionsToString(t, e.Clone()) + messageIn := tt.messageFactory(eventIn) + + err := WriteProducerMessage(tt.context, messageIn, record) + require.NoError(t, err) + + messageOut := NewMessage(record) + require.Equal(t, tt.expectedEncoding, messageOut.ReadEncoding()) + + if tt.expectedEncoding == binding.EncodingBinary { + err = messageOut.ReadBinary(tt.context, (*kafkaRecordWriter)(record)) + require.NoError(t, err) + } + + eventOut, err := binding.ToEvent(tt.context, messageOut) + require.NoError(t, err) + + if tt.expectedEncoding == binding.EncodingBinary { + eventIn.SetExtension(KafkaPartitionKey, strconv.FormatInt(int64(record.Partition), 10)) + eventIn.SetExtension(KafkaOffsetKey, strconv.FormatInt(record.Offset, 10)) + eventIn.SetExtension(KafkaTopicKey, record.Topic) + } + AssertEventEquals(t, eventIn, *eventOut) + }) + } + }) +}