This is an automated email from the ASF dual-hosted git repository.
jrmccluskey pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new b456c68217d [Go SDK] Use deterministic protobuf marshaling in protox
helpers (#40385)
b456c68217d is described below
commit b456c68217d13a079174b6c22d65fed721ac181a
Author: Malo Denielou <[email protected]>
AuthorDate: Tue Oct 6 14:32:31 2026 +0000
[Go SDK] Use deterministic protobuf marshaling in protox helpers (#40385)
---
sdks/go/pkg/beam/core/util/protox/any.go | 4 +--
sdks/go/pkg/beam/core/util/protox/any_test.go | 40 +++++++++++++++++++++++++++
sdks/go/pkg/beam/core/util/protox/base64.go | 2 +-
sdks/go/pkg/beam/core/util/protox/protox.go | 2 +-
4 files changed, 44 insertions(+), 4 deletions(-)
diff --git a/sdks/go/pkg/beam/core/util/protox/any.go
b/sdks/go/pkg/beam/core/util/protox/any.go
index 46bd08b1aff..9f21d9f5d90 100644
--- a/sdks/go/pkg/beam/core/util/protox/any.go
+++ b/sdks/go/pkg/beam/core/util/protox/any.go
@@ -45,7 +45,7 @@ func UnpackProto(data *protobuf.Any, ret proto.Message) error
{
// PackProto encodes a proto message and wraps into a BytesValue.
func PackProto(in proto.Message) (*protobuf.Any, error) {
- b, err := proto.Marshal(in)
+ b, err := proto.MarshalOptions{Deterministic: true}.Marshal(in)
if err != nil {
return nil, err
}
@@ -88,7 +88,7 @@ func UnpackBytes(data *protobuf.Any) ([]byte, error) {
func PackBytes(in []byte) (*protobuf.Any, error) {
var buf protobufw.BytesValue
buf.Value = in
- b, err := proto.Marshal(&buf)
+ b, err := proto.MarshalOptions{Deterministic: true}.Marshal(&buf)
if err != nil {
return nil, err
}
diff --git a/sdks/go/pkg/beam/core/util/protox/any_test.go
b/sdks/go/pkg/beam/core/util/protox/any_test.go
index 9eb7621db35..16c73afb63b 100644
--- a/sdks/go/pkg/beam/core/util/protox/any_test.go
+++ b/sdks/go/pkg/beam/core/util/protox/any_test.go
@@ -17,9 +17,11 @@ package protox
import (
"bytes"
+ "fmt"
"testing"
"google.golang.org/protobuf/proto"
+ "google.golang.org/protobuf/types/known/structpb"
protobufw "google.golang.org/protobuf/types/known/wrapperspb"
)
@@ -81,3 +83,41 @@ func TestBytesPackingInvertibility(t *testing.T) {
t.Errorf("Got %v, wanted %v", b, data)
}
}
+
+func TestDeterministicEncoding(t *testing.T) {
+ fields := make(map[string]*structpb.Value)
+ for i := 0; i < 32; i++ {
+ fields[fmt.Sprintf("key_%02d", i)] =
structpb.NewNumberValue(float64(i))
+ }
+ msg := &structpb.Struct{Fields: fields}
+
+ wantBytes := MustEncode(msg)
+ wantBase64, err := EncodeBase64(msg)
+ if err != nil {
+ t.Fatalf("EncodeBase64 failed: %v", err)
+ }
+ wantAny, err := PackProto(msg)
+ if err != nil {
+ t.Fatalf("PackProto failed: %v", err)
+ }
+
+ for i := 0; i < 50; i++ {
+ if got := MustEncode(msg); !bytes.Equal(got, wantBytes) {
+ t.Fatalf("MustEncode produced non-deterministic bytes
on iteration %d", i)
+ }
+ gotBase64, err := EncodeBase64(msg)
+ if err != nil {
+ t.Fatalf("EncodeBase64 failed: %v", err)
+ }
+ if gotBase64 != wantBase64 {
+ t.Fatalf("EncodeBase64 produced non-deterministic
output on iteration %d", i)
+ }
+ gotAny, err := PackProto(msg)
+ if err != nil {
+ t.Fatalf("PackProto failed: %v", err)
+ }
+ if !bytes.Equal(gotAny.GetValue(), wantAny.GetValue()) {
+ t.Fatalf("PackProto produced non-deterministic bytes on
iteration %d", i)
+ }
+ }
+}
diff --git a/sdks/go/pkg/beam/core/util/protox/base64.go
b/sdks/go/pkg/beam/core/util/protox/base64.go
index 79ea8a025f7..257c9213653 100644
--- a/sdks/go/pkg/beam/core/util/protox/base64.go
+++ b/sdks/go/pkg/beam/core/util/protox/base64.go
@@ -33,7 +33,7 @@ func MustEncodeBase64(msg proto.Message) string {
// EncodeBase64 encodes a proto wrapped in base64.
func EncodeBase64(msg proto.Message) (string, error) {
- data, err := proto.Marshal(msg)
+ data, err := proto.MarshalOptions{Deterministic: true}.Marshal(msg)
if err != nil {
return "", err
}
diff --git a/sdks/go/pkg/beam/core/util/protox/protox.go
b/sdks/go/pkg/beam/core/util/protox/protox.go
index 892a2ba97d0..d76fdc0d2f6 100644
--- a/sdks/go/pkg/beam/core/util/protox/protox.go
+++ b/sdks/go/pkg/beam/core/util/protox/protox.go
@@ -20,7 +20,7 @@ import "google.golang.org/protobuf/proto"
// MustEncode encode the message and panics on failure.
func MustEncode(msg proto.Message) []byte {
- data, err := proto.Marshal(msg)
+ data, err := proto.MarshalOptions{Deterministic: true}.Marshal(msg)
if err != nil {
panic(err)
}