From e1cd47e3f8f8e6cb877e6e18d26dc752030a2a3e Mon Sep 17 00:00:00 2001 From: sergeyb Date: Thu, 17 Sep 2026 22:40:41 +0000 Subject: [PATCH] test(consumer): wait for transport metrics before assertions MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Summary: Intent: - Remove a scheduling race that let consumer metrics tests snapshot before transport completion metrics were recorded. Changes: - Stop the consumer before taking metric snapshots so worker draining synchronizes Ack, Nack, and Reject completion. - Preserve the existing metric assertions without sleeps, retries, or weaker expectations. --- Generated by the 🪄 pr-create skill in devexp-agent-marketplace --- platform/consumer/consumer_test.go | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/platform/consumer/consumer_test.go b/platform/consumer/consumer_test.go index f7dd078f..ff3d68cf 100644 --- a/platform/consumer/consumer_test.go +++ b/platform/consumer/consumer_test.go @@ -849,6 +849,8 @@ func TestConsumer_ObservabilityTags(t *testing.T) { deliveryChan <- mockDel <-done + require.NoError(t, testC.Stop(30000)) + snapshot := testScope.Snapshot() histograms := snapshot.Histograms() @@ -895,8 +897,6 @@ func TestConsumer_ObservabilityTags(t *testing.T) { assert.NotContains(t, counter.Name(), duplicate) } } - - _ = testC.Stop(30000) }) } } @@ -989,6 +989,9 @@ func TestConsumer_AckLifecycleMetrics(t *testing.T) { deliveryChan <- mockDel <-done + err = c.Stop(30000) + require.NoError(t, err) + snapshot := scope.Snapshot() histograms := snapshot.Histograms() var foundAck bool @@ -999,9 +1002,6 @@ func TestConsumer_AckLifecycleMetrics(t *testing.T) { } } assert.True(t, foundAck, "Should have successful ack.finish metric") - - err = c.Stop(30000) - require.NoError(t, err) } func TestConsumer_NackLifecycleMetrics(t *testing.T) { @@ -1043,6 +1043,9 @@ func TestConsumer_NackLifecycleMetrics(t *testing.T) { deliveryChan <- mockDel <-done + err = c.Stop(30000) + require.NoError(t, err) + snapshot := scope.Snapshot() histograms := snapshot.Histograms() var foundNackError bool @@ -1053,9 +1056,6 @@ func TestConsumer_NackLifecycleMetrics(t *testing.T) { } } assert.True(t, foundNackError, "Should have failed nack.finish metric") - - err = c.Stop(30000) - require.NoError(t, err) } // TestConsumer_PerPartitionProcessing verifies that a slow message on partition A