@@ -47,9 +47,12 @@ class DeliveryLog:
4747 This class is used for verifying integrity of message flow from producer to consumers.
4848 """
4949
50- def __init__ (self , logger , allow_duplicates : bool ):
50+ def __init__ (
51+ self , logger , allow_duplicates : bool , allow_out_of_order : bool = False
52+ ):
5153 self ._logger = logger
5254 self ._allow_duplicates = allow_duplicates
55+ self ._allow_out_of_order = allow_out_of_order
5356
5457 self ._num_acks = 0
5558 self ._num_nacks = 0
@@ -283,15 +286,16 @@ def feed_consumer_log(self, path: Path, uri: str, app_id: str) -> None:
283286 assert self ._num_errors == 0 , self ._format_error_status ()
284287
285288 # Secondly, make sure that PUSHes are in correct order
286- previous_index = - 1
287- for message_index in self ._push_order :
288- if message_index < previous_index :
289- self ._error (
290- f"{ app_id } : out of order PUSH" ,
291- self ._pushes [message_index ],
292- self ._pushes [previous_index ],
293- )
294- previous_index = message_index
289+ if not self ._allow_out_of_order :
290+ previous_index = - 1
291+ for message_index in self ._push_order :
292+ if message_index < previous_index :
293+ self ._error (
294+ f"{ app_id } : out of order PUSH" ,
295+ self ._pushes [message_index ],
296+ self ._pushes [previous_index ],
297+ )
298+ previous_index = message_index
295299
296300 assert self ._num_errors == 0 , self ._format_error_status ()
297301 assert self ._consumed == self ._num_puts
@@ -332,7 +336,7 @@ class TestPutsRetransmission:
332336
333337 work_dir : Path
334338
335- def inspect_results (self , allow_duplicates = False ):
339+ def inspect_results (self , allow_duplicates = False , allow_out_of_order = False ):
336340 if self .active_node in self .cluster .virtual_nodes ():
337341 self .active_node .wait_status (wait_leader = True , wait_ready = False )
338342
@@ -357,10 +361,12 @@ def inspect_results(self, allow_duplicates=False):
357361 for consumer in self .consumers :
358362 consumer [0 ].force_stop ()
359363
360- self .parse_message_logs (allow_duplicates = allow_duplicates )
364+ self .parse_message_logs (
365+ allow_duplicates = allow_duplicates , allow_out_of_order = allow_out_of_order
366+ )
361367
362- def parse_message_logs (self , allow_duplicates = False ):
363- delivery_log = DeliveryLog (test_logger , allow_duplicates )
368+ def parse_message_logs (self , allow_duplicates = False , allow_out_of_order = False ):
369+ delivery_log = DeliveryLog (test_logger , allow_duplicates , allow_out_of_order )
364370 delivery_log .feed_producer_log (self .work_dir / "producer.log" , self .uri )
365371 for uri , app_id in self .uris :
366372 delivery_log .feed_consumer_log (self .work_dir / f"{ app_id } .log" , uri , app_id )
@@ -645,7 +651,7 @@ def test_kill_replica(self, multi_node: Cluster, domain_urls: tc.DomainUrls):
645651
646652 # Because the quorum is 3, cluster is still healthy after shutting down
647653 # replica.
648- self .inspect_results (allow_duplicates = True )
654+ self .inspect_results (allow_duplicates = True , allow_out_of_order = True )
649655
650656 @tweak .broker .app_config .network_interfaces .tcp_interface .low_watermark (512 )
651657 @tweak .broker .app_config .network_interfaces .tcp_interface .high_watermark (1024 )
@@ -672,6 +678,21 @@ def test_watermarks(self, multi_node: Cluster, domain_urls: tc.DomainUrls):
672678 self .inspect_results (allow_duplicates = False )
673679
674680 def test_kill_proxy (self , multi_node : Cluster , domain_urls : tc .DomainUrls ):
681+ """
682+ kill replica can result in out-of-order in the following scenario:
683+
684+ - consumer -> proxy -> replica1 -> primary.
685+ - kill replica1
686+ - proxy detects replica1 disconnect before primary does, and opens the
687+ queue on replica2
688+ - replica2 opens the queue on primary before primary detects replica1
689+ disconnect, resulting in primary having 2 downstreams.
690+ - primary round robins messages {m1, m2}. m1 goes to replica1, m2 goes
691+ to replica2
692+ - primary detects replica1 disconnect and redelivers m1 to replica2
693+ resulting in out-of-order {m2, m1}
694+ """
695+
675696 self .setup_cluster_fanout (multi_node , domain_urls )
676697
677698 self .replica_proxy .force_stop ()
0 commit comments