From 13dfc5fcdd0c1afefd99ff1e3dbb8417aa338c22 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Zolt=C3=A1n=20Papp?= Date: Tue, 31 Mar 2026 14:46:19 +0200 Subject: [PATCH] [client] Fix flow client Receive retry loop not stopping after Close MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Use backoff.Permanent for canceled gRPC errors so Receive returns immediately instead of retrying until context deadline when the connection is already closed. Add TestNewClient_PermanentClose to verify the behavior. The connectivity.Shutdown check was meaningless because when the connection is shut down, c.realClient.Events(ctx, grpc.WaitForReady(true)) on the nex line already fails with codes.Canceled — which is now handled as a permanent error. The explicit state check was just duplicating what gRPC already reports through its normal error path. --- flow/client/client.go | 42 ++++++++++++++++++-------------------- flow/client/client_test.go | 32 ++++++++++++++++++++++++----- 2 files changed, 47 insertions(+), 27 deletions(-) diff --git a/flow/client/client.go b/flow/client/client.go index 318fcfe1e..64aecb3fb 100644 --- a/flow/client/client.go +++ b/flow/client/client.go @@ -14,7 +14,6 @@ import ( log "github.com/sirupsen/logrus" "google.golang.org/grpc" "google.golang.org/grpc/codes" - "google.golang.org/grpc/connectivity" "google.golang.org/grpc/credentials" "google.golang.org/grpc/credentials/insecure" "google.golang.org/grpc/keepalive" @@ -88,12 +87,31 @@ func (c *GRPCClient) Close() error { return nil } +func (c *GRPCClient) Send(event *proto.FlowEvent) error { + c.streamMu.Lock() + stream := c.stream + c.streamMu.Unlock() + + if stream == nil { + return errors.New("stream not initialized") + } + + if err := stream.Send(event); err != nil { + return fmt.Errorf("send flow event: %w", err) + } + + return nil +} + func (c *GRPCClient) Receive(ctx context.Context, interval time.Duration, msgHandler func(msg *proto.FlowEventAck) error) error { backOff := defaultBackoff(ctx, interval) operation := func() error { if err := c.establishStreamAndReceive(ctx, msgHandler); err != nil { + if errors.Is(err, context.Canceled) { + return backoff.Permanent(err) + } if s, ok := status.FromError(err); ok && s.Code() == codes.Canceled { - return fmt.Errorf("receive: %w: %w", err, context.Canceled) + return backoff.Permanent(fmt.Errorf("receive: %w", err)) } log.Errorf("receive failed: %v", err) return fmt.Errorf("receive: %w", err) @@ -109,10 +127,6 @@ func (c *GRPCClient) Receive(ctx context.Context, interval time.Duration, msgHan } func (c *GRPCClient) establishStreamAndReceive(ctx context.Context, msgHandler func(msg *proto.FlowEventAck) error) error { - if c.clientConn.GetState() == connectivity.Shutdown { - return errors.New("connection to flow receiver has been shut down") - } - stream, err := c.realClient.Events(ctx, grpc.WaitForReady(true)) if err != nil { return fmt.Errorf("create event stream: %w", err) @@ -177,19 +191,3 @@ func defaultBackoff(ctx context.Context, interval time.Duration) backoff.BackOff Clock: backoff.SystemClock, }, ctx) } - -func (c *GRPCClient) Send(event *proto.FlowEvent) error { - c.streamMu.Lock() - stream := c.stream - c.streamMu.Unlock() - - if stream == nil { - return errors.New("stream not initialized") - } - - if err := stream.Send(event); err != nil { - return fmt.Errorf("send flow event: %w", err) - } - - return nil -} diff --git a/flow/client/client_test.go b/flow/client/client_test.go index efe01c003..0972c7b65 100644 --- a/flow/client/client_test.go +++ b/flow/client/client_test.go @@ -169,11 +169,6 @@ func TestReceive_ContextCancellation(t *testing.T) { assert.NoError(t, err, "failed to close flow") }) - go func() { - time.Sleep(100 * time.Millisecond) - cancel() - }() - handlerCalled := false msgHandler := func(msg *proto.FlowEventAck) error { if !msg.IsInitiator { @@ -254,3 +249,30 @@ func TestSend(t *testing.T) { t.Fatal("timeout waiting for ack to be received by flow") } } + +func TestNewClient_PermanentClose(t *testing.T) { + server := newTestServer(t) + + client, err := flow.NewClient("http://"+server.addr, "test-payload", "test-signature", 1*time.Second) + require.NoError(t, err) + + err = client.Close() + require.NoError(t, err) + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + t.Cleanup(cancel) + + done := make(chan error, 1) + go func() { + done <- client.Receive(ctx, 1*time.Second, func(msg *proto.FlowEventAck) error { + return nil + }) + }() + + select { + case err := <-done: + require.Error(t, err) + case <-time.After(2 * time.Second): + t.Fatal("Receive did not return after Close — stuck in retry loop") + } +}