From c8dc704e4df3b0fa644234b32afd55562ed89384 Mon Sep 17 00:00:00 2001 From: Michael McQuade Date: Thu, 1 Oct 2026 18:57:46 -0500 Subject: [PATCH] test: the dedup test holds the first delivery until the repeated publish --- messaging_test.go | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/messaging_test.go b/messaging_test.go index 1bd57c4..afe9948 100644 --- a/messaging_test.go +++ b/messaging_test.go @@ -133,8 +133,14 @@ func TestPublishDedupAndReplay(t *testing.T) { h := newHarness(t) ctx := context.Background() var calls sync.Map + // The first delivery is held until the repeated publish has been made: + // its dedup key is free again as soon as it is finalized. + release := make(chan struct{}) + unblock := sync.OnceFunc(func() { close(release) }) + defer unblock() workers := hopper.NewWorkers() hopper.Subscribe(workers, hopper.Subscription{Name: "flaky", Pattern: "allocation.created", MaxAttempts: 1}, func(_ context.Context, msg *hopper.Message[allocationCreated]) error { + <-release n, _ := calls.LoadOrStore(msg.ID, new(int)) *n.(*int)++ if *n.(*int) == 1 { @@ -155,6 +161,7 @@ func TestPublishDedupAndReplay(t *testing.T) { if len(again.Deliveries) != 1 || !again.Deliveries[0].Duplicate || again.Deliveries[0].Job.ID != first.Deliveries[0].Job.ID { t.Errorf("repeated publish = %+v", again.Deliveries[0]) } + unblock() // The single attempt fails and the delivery is dead-lettered; replay // re-drives it and it succeeds.