diff --git a/pulsar-function-go/go.mod b/pulsar-function-go/go.mod index 576577dea1048..0c08f4129fb46 100644 --- a/pulsar-function-go/go.mod +++ b/pulsar-function-go/go.mod @@ -4,7 +4,6 @@ go 1.25.0 require ( github.com/apache/pulsar-client-go v0.20.0 - github.com/golang/protobuf v1.5.4 github.com/prometheus/client_golang v1.20.5 github.com/prometheus/client_model v0.6.1 github.com/sirupsen/logrus v1.9.3 diff --git a/pulsar-function-go/pf/instance.go b/pulsar-function-go/pf/instance.go index 8624120d57801..f2b5dfd50b801 100644 --- a/pulsar-function-go/pf/instance.go +++ b/pulsar-function-go/pf/instance.go @@ -27,13 +27,11 @@ import ( "strings" "time" - "github.com/golang/protobuf/ptypes/empty" - "github.com/apache/pulsar-client-go/pulsar" - log "github.com/apache/pulsar/pulsar-function-go/logutil" pb "github.com/apache/pulsar/pulsar-function-go/pb" prometheus_client "github.com/prometheus/client_model/go" + "google.golang.org/protobuf/types/known/emptypb" ) type goInstance struct { @@ -648,9 +646,9 @@ func (gi *goInstance) getAndResetMetrics() *pb.MetricsData { return metricsData } -func (gi *goInstance) resetMetrics() *empty.Empty { +func (gi *goInstance) resetMetrics() *emptypb.Empty { gi.stats.reset() - return &empty.Empty{} + return &emptypb.Empty{} } // This method is used to get the required metrics for Prometheus. diff --git a/pulsar-function-go/pf/stats_test.go b/pulsar-function-go/pf/stats_test.go index 138dc91cd9cd3..550ef18aa8f0e 100644 --- a/pulsar-function-go/pf/stats_test.go +++ b/pulsar-function-go/pf/stats_test.go @@ -28,10 +28,10 @@ import ( "testing" "time" - "github.com/golang/protobuf/ptypes/empty" "github.com/prometheus/client_golang/prometheus" "github.com/stretchr/testify/assert" "google.golang.org/protobuf/encoding/prototext" + "google.golang.org/protobuf/types/known/emptypb" prometheus_client "github.com/prometheus/client_model/go" ) @@ -196,7 +196,7 @@ func TestInstanceControlMetrics(t *testing.T) { instance := newGoInstance() t.Cleanup(instance.close) instanceClient := instanceCommunicationClient(t, instance) - _, err := instanceClient.GetMetrics(context.Background(), &empty.Empty{}) + _, err := instanceClient.GetMetrics(context.Background(), &emptypb.Empty{}) assert.NoError(t, err, "err communicating with instance control: %v", err) testLabels := []string{"userMetricControlTest1", "userMetricControlTest2"} @@ -209,7 +209,7 @@ func TestInstanceControlMetrics(t *testing.T) { } time.Sleep(time.Second) - metrics, err := instanceClient.GetMetrics(context.Background(), &empty.Empty{}) + metrics, err := instanceClient.GetMetrics(context.Background(), &emptypb.Empty{}) assert.NoError(t, err, "err communicating with instance control: %v", err) for value, label := range testLabels { assert.Containsf(t, metrics.UserMetrics, label, "user metrics should contain metric %s", label)