Fix IntReactiveUtils for the proper emission

* Enable `MessageChannelReactiveUtilsTests.testOverproducingWithSubscribableChannel()` back
* Disable `StompServerIntegrationTests` again since build on CI is stalled again
This commit is contained in:
Artem Bilan
2020-08-10 12:32:47 -04:00
parent d14c5d5841
commit 17fc4eae99
3 changed files with 8 additions and 4 deletions

View File

@@ -17,6 +17,7 @@
package org.springframework.integration.util;
import java.time.Duration;
import java.util.concurrent.locks.LockSupport;
import org.reactivestreams.Publisher;
@@ -127,7 +128,11 @@ public final class IntegrationReactiveUtils {
return Flux.defer(() -> {
Sinks.Many<Message<T>> sink = Sinks.many().multicast().onBackpressureBuffer(1);
@SuppressWarnings("unchecked")
MessageHandler messageHandler = (message) -> sink.emitNext((Message<T>) message);
MessageHandler messageHandler = (message) -> {
while (!sink.emitNext((Message<T>) message).hasEmitted()) {
LockSupport.parkNanos(10);
}
};
inputChannel.subscribe(messageHandler);
return sink.asFlux()
.doOnCancel(() -> inputChannel.unsubscribe(messageHandler));

View File

@@ -21,7 +21,6 @@ import static org.assertj.core.api.Assertions.assertThat;
import java.time.Duration;
import java.util.concurrent.atomic.AtomicInteger;
import org.junit.jupiter.api.Disabled;
import org.junit.jupiter.api.Test;
import org.springframework.integration.util.IntegrationReactiveUtils;
@@ -71,7 +70,6 @@ class MessageChannelReactiveUtilsTests {
}
@Test
@Disabled("Backpressure is not honored")
void testOverproducingWithSubscribableChannel() {
DirectChannel channel = new DirectChannel();

View File

@@ -22,6 +22,7 @@ import static org.assertj.core.api.Assertions.assertThatExceptionOfType;
import org.apache.activemq.broker.BrokerService;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Disabled;
import org.junit.jupiter.api.Test;
import org.springframework.context.ApplicationEvent;
@@ -62,7 +63,7 @@ import org.springframework.util.SocketUtils;
*
* @since 4.2
*/
//@Disabled("Until the fix in reactor-netty-core")
@Disabled("Until the fix in reactor-netty-core")
public class StompServerIntegrationTests {
private static BrokerService activeMQBroker;