Fix compatibility with latest Spring AMQP
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2014-2022 the original author or authors.
|
||||
* Copyright 2014-2023 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -20,6 +20,7 @@ import java.util.Collection;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.CyclicBarrier;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.locks.Lock;
|
||||
|
||||
import org.junit.After;
|
||||
import org.junit.ClassRule;
|
||||
@@ -134,19 +135,23 @@ public class ChannelTests {
|
||||
private void waitForNewConsumer(PublishSubscribeAmqpChannel channel, BlockingQueueConsumer consumer)
|
||||
throws Exception {
|
||||
|
||||
final Object consumersMonitor = TestUtils.getPropertyValue(channel, "container.consumersMonitor");
|
||||
Lock consumersLock = TestUtils.getPropertyValue(channel, "container.consumersLock", Lock.class);
|
||||
int n = 0;
|
||||
while (n++ < 100) {
|
||||
Set<BlockingQueueConsumer> consumers = TestUtils
|
||||
.getPropertyValue(channel, "container.consumers", Set.class);
|
||||
synchronized (consumersMonitor) {
|
||||
consumersLock.lock();
|
||||
try {
|
||||
if (!consumers.isEmpty()) {
|
||||
BlockingQueueConsumer newConsumer = consumers.iterator().next();
|
||||
if (newConsumer != consumer && newConsumer.getConsumerTags().size() > 0) {
|
||||
if (newConsumer != consumer && !newConsumer.getConsumerTags().isEmpty()) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
finally {
|
||||
consumersLock.unlock();
|
||||
}
|
||||
Thread.sleep(100);
|
||||
}
|
||||
assertThat(n < 100).as("Failed to restart consumer").isTrue();
|
||||
|
||||
Reference in New Issue
Block a user