@@ -147,16 +147,16 @@ async def publish_outbound_side_effect_logic(*args, **kwargs):
147147 ):
148148 with attempt :
149149 assert_that (len (captured_logs )).is_greater_than (0 )
150- retry_consumer_logs = list (
150+ retry_sent_to_inbound_logs = list (
151151 filter (
152152 lambda entry : str (entry ["event" ]).startswith (
153153 "Sent to inbound queue"
154154 ),
155155 captured_logs ,
156156 )
157157 )
158- assert_that (retry_consumer_logs ).is_length (2 )
159- retry_headers : HeadersType = retry_consumer_logs [0 ]["headers" ]
158+ assert_that (retry_sent_to_inbound_logs ).is_length (2 )
159+ retry_headers : HeadersType = retry_sent_to_inbound_logs [0 ]["headers" ]
160160 assert_that (retry_headers .get ("test" )).is_equal_to (headers .get ("test" ))
161161 assert_that (retry_headers .get ("message_id" )).is_equal_to (
162162 headers .get ("message_id" )
@@ -170,7 +170,7 @@ async def publish_outbound_side_effect_logic(*args, **kwargs):
170170 datetime .now (UTC ) - timedelta (seconds = 10 ),
171171 datetime .now (UTC ) + timedelta (seconds = 10 ),
172172 )
173- retry_payload_string = retry_consumer_logs [0 ]["payload_string" ]
173+ retry_payload_string = retry_sent_to_inbound_logs [0 ]["payload_string" ]
174174 outbound_samples : list [dict [str , Any ]] = json .loads (retry_payload_string )
175175 assert_that (
176176 DeepDiff (
@@ -180,6 +180,14 @@ async def publish_outbound_side_effect_logic(*args, **kwargs):
180180 )
181181 ).is_empty ()
182182
183+ retry_sent_to_retry_logs = list (
184+ filter (
185+ lambda entry : str (entry ["event" ]).startswith ("Sent to retry queue" ),
186+ captured_logs ,
187+ )
188+ )
189+ assert_that (len (retry_sent_to_retry_logs )).is_greater_than (1 )
190+
183191 for attempt in Retrying (
184192 stop = stop_after_delay (5 ), wait = wait_fixed (0.5 ), reraise = True
185193 ):
0 commit comments