8282import org .junit .jupiter .params .provider .CsvSource ;
8383import org .junit .jupiter .params .provider .EnumSource ;
8484import org .junit .jupiter .params .provider .ValueSource ;
85+ import org .slf4j .Logger ;
86+ import org .slf4j .LoggerFactory ;
8587
8688@ AmqpTestInfrastructure
8789public class AmqpTest {
8890
91+ private static final Logger LOGGER = LoggerFactory .getLogger (AmqpTest .class );
92+
8993 Connection connection ;
9094 Environment environment ;
9195 String name ;
@@ -1636,6 +1640,12 @@ void asyncMessageAcceptanceWithExecutorService() {
16361640
16371641 AtomicReference <ExecutorService > executorService = new AtomicReference <>();
16381642
1643+ Sync consumeSync = sync (messageCount );
1644+ AtomicInteger received = new AtomicInteger ();
1645+ AtomicInteger dispatched = new AtomicInteger ();
1646+ AtomicInteger executed = new AtomicInteger ();
1647+ AtomicInteger processed = new AtomicInteger ();
1648+ AtomicInteger accepted = new AtomicInteger ();
16391649 try {
16401650 Sync publishSync = sync (messageCount );
16411651 range (0 , messageCount )
@@ -1646,31 +1656,35 @@ void asyncMessageAcceptanceWithExecutorService() {
16461656
16471657 assertThat (publishSync ).completes ();
16481658 publisher .close ();
1659+ LOGGER .debug ("Published {} messages in queue '{}'" , messageCount , name );
16491660
16501661 executorService .set (Executors .newFixedThreadPool (executorThreads ));
1651- Sync consumeSync = sync (messageCount );
16521662 long [] simulatedLatencies = new long [executorThreads ];
16531663 Random random = new Random ();
16541664 for (int i = 0 ; i < executorThreads ; i ++) {
16551665 simulatedLatencies [i ] = random .nextInt (4 ) + 1 ;
16561666 }
16571667 AtomicInteger count = new AtomicInteger ();
16581668 Supplier <Long > latency = () -> simulatedLatencies [count .getAndIncrement () % executorThreads ];
1659-
16601669 consumer =
16611670 connection
16621671 .consumerBuilder ()
16631672 .queue (name )
16641673 .messageHandler (
16651674 (context , message ) -> {
1675+ received .incrementAndGet ();
16661676 executorService
16671677 .get ()
16681678 .submit (
16691679 () -> {
1680+ executed .incrementAndGet ();
16701681 TestUtils .simulateActivity (latency .get ());
1682+ processed .incrementAndGet ();
16711683 context .accept ();
1684+ accepted .incrementAndGet ();
16721685 consumeSync .down ();
16731686 });
1687+ dispatched .incrementAndGet ();
16741688 })
16751689 .build ();
16761690
@@ -1679,6 +1693,19 @@ void asyncMessageAcceptanceWithExecutorService() {
16791693 Management .QueueInfo queueInfo = connection .management ().queueInfo (name );
16801694 assertThat (queueInfo ).hasName (name ).hasMessageCount (0 );
16811695 } finally {
1696+ if (consumer != null ) {
1697+ LOGGER .debug ("Consumer state: '{}'" , ((AmqpConsumer ) consumer ).diagnosticState ());
1698+ }
1699+ LOGGER .debug (
1700+ "Counters: received={}, dispatched={}, executed={}, processed={}, accepted={}" ,
1701+ received .get (),
1702+ dispatched .get (),
1703+ executed .get (),
1704+ processed .get (),
1705+ accepted .get ());
1706+ LOGGER .debug ("Consume latch: {}" , consumeSync .count ());
1707+ LOGGER .debug ("Queue info: {}" , connection .management ().queueInfo (name ));
1708+
16821709 if (consumer != null ) {
16831710 consumer .close ();
16841711 }
0 commit comments