From 367bee9aeda972f9613ce82a8c2cc7b33e8e9e49 Mon Sep 17 00:00:00 2001 From: Jem Davies Date: Mon, 1 Jun 2026 12:52:10 +0100 Subject: [PATCH 1/6] add goleak to protobuf processor test Signed-off-by: Jem Davies --- go.mod | 1 + .../impl/protobuf/processor_protobuf_test.go | 41 ++++++++++++++----- 2 files changed, 32 insertions(+), 10 deletions(-) diff --git a/go.mod b/go.mod index 074f72a5a5..0d2d20a516 100644 --- a/go.mod +++ b/go.mod @@ -157,6 +157,7 @@ require ( go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.40.0 go.opentelemetry.io/otel/sdk v1.43.0 go.opentelemetry.io/otel/trace v1.43.0 + go.uber.org/goleak v1.3.0 go.uber.org/multierr v1.11.0 golang.org/x/crypto v0.52.0 golang.org/x/net v0.55.0 diff --git a/internal/impl/protobuf/processor_protobuf_test.go b/internal/impl/protobuf/processor_protobuf_test.go index 7b2f0cc4ce..a2b128e476 100644 --- a/internal/impl/protobuf/processor_protobuf_test.go +++ b/internal/impl/protobuf/processor_protobuf_test.go @@ -2,7 +2,6 @@ package protobuf import ( "context" - "errors" "fmt" "net" "net/http" @@ -14,6 +13,7 @@ import ( "connectrpc.com/connect" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "go.uber.org/goleak" "google.golang.org/protobuf/reflect/protodesc" "google.golang.org/protobuf/reflect/protoreflect" "google.golang.org/protobuf/reflect/protoregistry" @@ -22,6 +22,10 @@ import ( "github.com/warpstreamlabs/bento/public/service" ) +func TestMain(m *testing.M) { + goleak.VerifyTestMain(m) +} + const protosPath = "../../../config/test/protobuf/schema" func TestProtobufFromJSON(t *testing.T) { @@ -103,6 +107,7 @@ discard_unknown: %t assert.Contains(t, string(mBytes), exp) } require.NoError(t, msgs[0].GetError()) + }) t.Run(test.name+" bsr", func(t *testing.T) { @@ -125,6 +130,9 @@ discard_unknown: %t require.NoError(t, res) require.Len(t, msgs, 1) + err = proc.Close(context.Background()) + require.NoError(t, err) + mBytes, err := msgs[0].AsBytes() require.NoError(t, err) @@ -238,6 +246,9 @@ emit_unpopulated: %t require.NoError(t, res) require.Len(t, msgs, 1) + err = proc.Close(context.Background()) + require.NoError(t, err) + mBytes, err := msgs[0].AsBytes() require.NoError(t, err) @@ -266,6 +277,9 @@ emit_unpopulated: %t require.NoError(t, res) require.Len(t, msgs, 1) + err = proc.Close(context.Background()) + require.NoError(t, err) + mBytes, err := msgs[0].AsBytes() require.NoError(t, err) @@ -324,6 +338,9 @@ import_paths: [ %v ] _, err = proc.Process(context.Background(), service.NewMessage([]byte(test.input))) require.Error(t, err) require.Contains(t, err.Error(), test.output) + + err = proc.Close(context.Background()) + require.NoError(t, err) }) t.Run(test.name+" bsr", func(tt *testing.T) { @@ -344,6 +361,9 @@ bsr: _, err = proc.Process(context.Background(), service.NewMessage([]byte(test.input))) require.Error(t, err) require.Contains(t, err.Error(), test.output) + + err = proc.Close(context.Background()) + require.NoError(t, err) }) } } @@ -456,16 +476,17 @@ func runMockBSRServer(t *testing.T) string { mux := http.NewServeMux() fileDescriptorSetServer := &fileDescriptorSetServer{fileDescriptorSet: fileDescriptorSet} mux.Handle(reflectv1beta1connect.NewFileDescriptorSetServiceHandler(fileDescriptorSetServer)) - go func() { - srv := &http.Server{Handler: mux} - srv.Protocols = new(http.Protocols) - srv.Protocols.SetHTTP1(true) - srv.Protocols.SetUnencryptedHTTP2(true) - if err := http.Serve(listener, srv.Handler); err != nil && !errors.Is(err, http.ErrServerClosed) { - require.NoError(t, err) - } - }() + srv := &http.Server{Handler: mux} + srv.Protocols = new(http.Protocols) + srv.Protocols.SetHTTP1(true) + srv.Protocols.SetUnencryptedHTTP2(true) + + go srv.Serve(listener) + + t.Cleanup(func() { + srv.Close() + }) return listener.Addr().String() } From 5396bc77f5e98bea82d30c063371cbadffe5ecc2 Mon Sep 17 00:00:00 2001 From: Jem Davies Date: Mon, 1 Jun 2026 12:55:31 +0100 Subject: [PATCH 2/6] remove go-routine leaks Signed-off-by: Jem Davies --- internal/impl/protobuf/multimodule_watcher.go | 30 +++++++++++++++---- internal/impl/protobuf/processor_protobuf.go | 6 +++- 2 files changed, 30 insertions(+), 6 deletions(-) diff --git a/internal/impl/protobuf/multimodule_watcher.go b/internal/impl/protobuf/multimodule_watcher.go index 463ea7a269..13ec58cdd2 100644 --- a/internal/impl/protobuf/multimodule_watcher.go +++ b/internal/impl/protobuf/multimodule_watcher.go @@ -19,6 +19,8 @@ import ( type MultiModuleWatcher struct { bsrClients map[string]*prototransform.SchemaWatcher + httpClient *http.Client + cancel context.CancelFunc } var _ prototransform.Resolver = &MultiModuleWatcher{} @@ -27,10 +29,18 @@ func newMultiModuleWatcher(bsrModules []*service.ParsedConfig) (*MultiModuleWatc if len(bsrModules) == 0 { return nil, errors.New("no modules provided") } - multiModuleWatcher := &MultiModuleWatcher{} + + ctx, cancel := context.WithCancel(context.Background()) + httpClient := &http.Client{ + Transport: &http.Transport{}, + } + multiModuleWatcher := &MultiModuleWatcher{ + bsrClients: make(map[string]*prototransform.SchemaWatcher), + httpClient: httpClient, + cancel: cancel, + } // Initialise one client for each module - multiModuleWatcher.bsrClients = make(map[string]*prototransform.SchemaWatcher) for _, bsrModule := range bsrModules { var bsrURL string bsrURL, err := bsrModule.FieldString(fieldBSRUrl) @@ -53,7 +63,7 @@ func newMultiModuleWatcher(bsrModules []*service.ParsedConfig) (*MultiModuleWatc return nil, err } - watcher, err := newSchemaWatcher(context.Background(), bsrURL, bsrAPIKey, module, version) + watcher, err := newSchemaWatcher(ctx, httpClient, bsrURL, bsrAPIKey, module, version) if err != nil { return nil, err } @@ -63,7 +73,7 @@ func newMultiModuleWatcher(bsrModules []*service.ParsedConfig) (*MultiModuleWatc return multiModuleWatcher, nil } -func newSchemaWatcher(ctx context.Context, bsrURL string, bsrAPIKey string, module string, version string) (*prototransform.SchemaWatcher, error) { +func newSchemaWatcher(ctx context.Context, httpClient *http.Client, bsrURL string, bsrAPIKey string, module string, version string) (*prototransform.SchemaWatcher, error) { // If no BSR url provided, extract from module if bsrURL == "" { segments := strings.Split(module, "/") @@ -80,7 +90,7 @@ func newSchemaWatcher(ctx context.Context, bsrURL string, bsrAPIKey string, modu if bsrAPIKey != "" { opts = append(opts, connectrpc.WithInterceptors(prototransform.NewAuthInterceptor(bsrAPIKey))) } - client := reflectv1beta1connect.NewFileDescriptorSetServiceClient(http.DefaultClient, bsrURL, opts...) + client := reflectv1beta1connect.NewFileDescriptorSetServiceClient(httpClient, bsrURL, opts...) cfg := &prototransform.SchemaWatcherConfig{ SchemaPoller: prototransform.NewSchemaPoller( @@ -173,3 +183,13 @@ func (w *MultiModuleWatcher) FindEnumByName(enum protoreflect.FullName) (protore } return nil, fmt.Errorf("could not find %s in any loaded modules", enum) } + +func (m *MultiModuleWatcher) Close() { + m.cancel() + for _, v := range m.bsrClients { + v.Stop() + } + if t, ok := m.httpClient.Transport.(*http.Transport); ok { + t.CloseIdleConnections() + } +} diff --git a/internal/impl/protobuf/processor_protobuf.go b/internal/impl/protobuf/processor_protobuf.go index f6f028a526..cf30863ee1 100644 --- a/internal/impl/protobuf/processor_protobuf.go +++ b/internal/impl/protobuf/processor_protobuf.go @@ -549,6 +549,10 @@ func (p *protobufProc) Process(ctx context.Context, msg *service.Message) (servi return service.MessageBatch{msg}, nil } -func (p *protobufProc) Close(context.Context) error { +func (p *protobufProc) Close(ctx context.Context) error { + if p.multiModuleWatcher != nil { + fmt.Println("Shutting Down watcher") + p.multiModuleWatcher.Close() + } return nil } From 78c958119354be723ce0eddc25ced25735f2ade7 Mon Sep 17 00:00:00 2001 From: Jem Davies <131159520+jem-davies@users.noreply.github.com> Date: Mon, 1 Jun 2026 12:57:45 +0100 Subject: [PATCH 3/6] Update internal/impl/protobuf/processor_protobuf.go --- internal/impl/protobuf/processor_protobuf.go | 1 - 1 file changed, 1 deletion(-) diff --git a/internal/impl/protobuf/processor_protobuf.go b/internal/impl/protobuf/processor_protobuf.go index cf30863ee1..dc4d9288b6 100644 --- a/internal/impl/protobuf/processor_protobuf.go +++ b/internal/impl/protobuf/processor_protobuf.go @@ -551,7 +551,6 @@ func (p *protobufProc) Process(ctx context.Context, msg *service.Message) (servi func (p *protobufProc) Close(ctx context.Context) error { if p.multiModuleWatcher != nil { - fmt.Println("Shutting Down watcher") p.multiModuleWatcher.Close() } return nil From af20c3ddde8bdbf5b67c3c7c7857ff0e9263226a Mon Sep 17 00:00:00 2001 From: Jem Davies Date: Mon, 1 Jun 2026 13:10:19 +0100 Subject: [PATCH 4/6] fix lint Signed-off-by: Jem Davies --- internal/impl/protobuf/processor_protobuf_test.go | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/internal/impl/protobuf/processor_protobuf_test.go b/internal/impl/protobuf/processor_protobuf_test.go index a2b128e476..e6da6a387a 100644 --- a/internal/impl/protobuf/processor_protobuf_test.go +++ b/internal/impl/protobuf/processor_protobuf_test.go @@ -2,6 +2,7 @@ package protobuf import ( "context" + "errors" "fmt" "net" "net/http" @@ -482,7 +483,11 @@ func runMockBSRServer(t *testing.T) string { srv.Protocols.SetHTTP1(true) srv.Protocols.SetUnencryptedHTTP2(true) - go srv.Serve(listener) + go func() { + if err := srv.Serve(listener); err != nil && !errors.Is(err, http.ErrServerClosed) { + t.Errorf("mock BSR server error: %v", err) + } + }() t.Cleanup(func() { srv.Close() From 2821b5014a11ab30d745eeb8c23682e7221299a0 Mon Sep 17 00:00:00 2001 From: Jem Davies Date: Tue, 2 Jun 2026 14:26:53 +0100 Subject: [PATCH 5/6] use a shutdown signaller instead of stored ctx Signed-off-by: Jem Davies --- internal/impl/protobuf/multimodule_watcher.go | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/internal/impl/protobuf/multimodule_watcher.go b/internal/impl/protobuf/multimodule_watcher.go index 13ec58cdd2..40cd20211c 100644 --- a/internal/impl/protobuf/multimodule_watcher.go +++ b/internal/impl/protobuf/multimodule_watcher.go @@ -10,6 +10,7 @@ import ( "buf.build/gen/go/bufbuild/reflect/connectrpc/go/buf/reflect/v1beta1/reflectv1beta1connect" connectrpc "connectrpc.com/connect" + "github.com/Jeffail/shutdown" "github.com/bufbuild/prototransform" "google.golang.org/protobuf/reflect/protoreflect" "google.golang.org/protobuf/reflect/protoregistry" @@ -20,7 +21,7 @@ import ( type MultiModuleWatcher struct { bsrClients map[string]*prototransform.SchemaWatcher httpClient *http.Client - cancel context.CancelFunc + shutSig *shutdown.Signaller } var _ prototransform.Resolver = &MultiModuleWatcher{} @@ -30,16 +31,17 @@ func newMultiModuleWatcher(bsrModules []*service.ParsedConfig) (*MultiModuleWatc return nil, errors.New("no modules provided") } - ctx, cancel := context.WithCancel(context.Background()) httpClient := &http.Client{ Transport: &http.Transport{}, } multiModuleWatcher := &MultiModuleWatcher{ bsrClients: make(map[string]*prototransform.SchemaWatcher), httpClient: httpClient, - cancel: cancel, + shutSig: shutdown.NewSignaller(), } + ctx, _ := multiModuleWatcher.shutSig.SoftStopCtx(context.Background()) + // Initialise one client for each module for _, bsrModule := range bsrModules { var bsrURL string @@ -185,7 +187,7 @@ func (w *MultiModuleWatcher) FindEnumByName(enum protoreflect.FullName) (protore } func (m *MultiModuleWatcher) Close() { - m.cancel() + m.shutSig.TriggerHardStop() for _, v := range m.bsrClients { v.Stop() } From 383d9361238fd87e8a27a27fcbb95a1ab7b35004 Mon Sep 17 00:00:00 2001 From: Jem Davies Date: Tue, 2 Jun 2026 14:32:59 +0100 Subject: [PATCH 6/6] clear map values on Close() calls Signed-off-by: Jem Davies --- internal/impl/protobuf/multimodule_watcher.go | 1 + internal/impl/protobuf/processor_protobuf.go | 1 + 2 files changed, 2 insertions(+) diff --git a/internal/impl/protobuf/multimodule_watcher.go b/internal/impl/protobuf/multimodule_watcher.go index 40cd20211c..5cc6a025c7 100644 --- a/internal/impl/protobuf/multimodule_watcher.go +++ b/internal/impl/protobuf/multimodule_watcher.go @@ -191,6 +191,7 @@ func (m *MultiModuleWatcher) Close() { for _, v := range m.bsrClients { v.Stop() } + m.bsrClients = nil if t, ok := m.httpClient.Transport.(*http.Transport); ok { t.CloseIdleConnections() } diff --git a/internal/impl/protobuf/processor_protobuf.go b/internal/impl/protobuf/processor_protobuf.go index dc4d9288b6..eddc02e291 100644 --- a/internal/impl/protobuf/processor_protobuf.go +++ b/internal/impl/protobuf/processor_protobuf.go @@ -552,6 +552,7 @@ func (p *protobufProc) Process(ctx context.Context, msg *service.Message) (servi func (p *protobufProc) Close(ctx context.Context) error { if p.multiModuleWatcher != nil { p.multiModuleWatcher.Close() + p.multiModuleWatcher = nil } return nil }