diff --git a/rsocket-core/src/jmh/java/io/rsocket/FragmentationPerf.java b/rsocket-core/src/jmh/java/io/rsocket/FragmentationPerf.java index 9d1510682..c5a76b09c 100644 --- a/rsocket-core/src/jmh/java/io/rsocket/FragmentationPerf.java +++ b/rsocket-core/src/jmh/java/io/rsocket/FragmentationPerf.java @@ -43,13 +43,15 @@ public void setup(Blackhole bh) { ByteBuffer data = createRandomBytes(1 << 18); ByteBuffer metadata = createRandomBytes(1 << 18); largeFrame = - Frame.Request.from(1, FrameType.REQUEST_RESPONSE, DefaultPayload.create(data, metadata), 1); + Frame.Request.from( + 1, FrameType.REQUEST_RESPONSE, DefaultPayload.create(data, metadata), 1); largeFrameFragmenter = new FrameFragmenter(1024); data = createRandomBytes(16); metadata = createRandomBytes(16); smallFrame = - Frame.Request.from(1, FrameType.REQUEST_RESPONSE, DefaultPayload.create(data, metadata), 1); + Frame.Request.from( + 1, FrameType.REQUEST_RESPONSE, DefaultPayload.create(data, metadata), 1); smallFrameFragmenter = new FrameFragmenter(2); smallFramesIterable = smallFrameFragmenter diff --git a/rsocket-core/src/jmh/java/io/rsocket/RSocketPerf.java b/rsocket-core/src/jmh/java/io/rsocket/RSocketPerf.java index 39b8bf44d..2391628c5 100644 --- a/rsocket-core/src/jmh/java/io/rsocket/RSocketPerf.java +++ b/rsocket-core/src/jmh/java/io/rsocket/RSocketPerf.java @@ -19,7 +19,6 @@ import io.rsocket.RSocketFactory.Start; import io.rsocket.perfutil.TestDuplexConnection; import io.rsocket.util.DefaultPayload; - import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; import org.openjdk.jmh.annotations.Benchmark; diff --git a/rsocket-core/src/main/java/io/rsocket/ConnectionSetupPayload.java b/rsocket-core/src/main/java/io/rsocket/ConnectionSetupPayload.java index 35feedbcf..49724beff 100644 --- a/rsocket-core/src/main/java/io/rsocket/ConnectionSetupPayload.java +++ b/rsocket-core/src/main/java/io/rsocket/ConnectionSetupPayload.java @@ -95,6 +95,7 @@ public ConnectionSetupPayload retain(int increment) { } public abstract ConnectionSetupPayload touch(); + public abstract ConnectionSetupPayload touch(Object hint); private static final class DefaultConnectionSetupPayload extends ConnectionSetupPayload { @@ -106,11 +107,7 @@ private static final class DefaultConnectionSetupPayload extends ConnectionSetup private final int flags; public DefaultConnectionSetupPayload( - String metadataMimeType, - String dataMimeType, - ByteBuf data, - ByteBuf metadata, - int flags) { + String metadataMimeType, String dataMimeType, ByteBuf data, ByteBuf metadata, int flags) { this.metadataMimeType = metadataMimeType; this.dataMimeType = dataMimeType; this.data = data; diff --git a/rsocket-core/src/main/java/io/rsocket/Frame.java b/rsocket-core/src/main/java/io/rsocket/Frame.java index 7dc8572e2..dbe4f05d8 100644 --- a/rsocket-core/src/main/java/io/rsocket/Frame.java +++ b/rsocket-core/src/main/java/io/rsocket/Frame.java @@ -273,7 +273,8 @@ public static Frame from( String metadataMimeType, String dataMimeType, Payload payload) { - final ByteBuf metadata = payload.hasMetadata() ? payload.sliceMetadata() : Unpooled.EMPTY_BUFFER; + final ByteBuf metadata = + payload.hasMetadata() ? payload.sliceMetadata() : Unpooled.EMPTY_BUFFER; final ByteBuf data = payload.sliceData(); final Frame frame = RECYCLER.get(); diff --git a/rsocket-core/src/main/java/io/rsocket/Payload.java b/rsocket-core/src/main/java/io/rsocket/Payload.java index 0b8f5ee8a..a00fe2e76 100644 --- a/rsocket-core/src/main/java/io/rsocket/Payload.java +++ b/rsocket-core/src/main/java/io/rsocket/Payload.java @@ -16,12 +16,10 @@ package io.rsocket; import io.netty.buffer.ByteBuf; - -import java.nio.ByteBuffer; -import java.nio.charset.StandardCharsets; - import io.netty.util.ReferenceCounted; import io.netty.util.ResourceLeakDetector; +import java.nio.ByteBuffer; +import java.nio.charset.StandardCharsets; /** Payload of a {@link Frame}. */ public interface Payload extends ReferenceCounted { @@ -47,30 +45,26 @@ public interface Payload extends ReferenceCounted { */ ByteBuf sliceData(); - /** - * Increases the reference count by {@code 1}. - */ + /** Increases the reference count by {@code 1}. */ @Override Payload retain(); - /** - * Increases the reference count by the specified {@code increment}. - */ + /** Increases the reference count by the specified {@code increment}. */ @Override Payload retain(int increment); /** - * Records the current access location of this object for debugging purposes. - * If this object is determined to be leaked, the information recorded by this operation will be provided to you - * via {@link ResourceLeakDetector}. This method is a shortcut to {@link #touch(Object) touch(null)}. + * Records the current access location of this object for debugging purposes. If this object is + * determined to be leaked, the information recorded by this operation will be provided to you via + * {@link ResourceLeakDetector}. This method is a shortcut to {@link #touch(Object) touch(null)}. */ @Override Payload touch(); /** - * Records the current access location of this object with an additional arbitrary information for debugging - * purposes. If this object is determined to be leaked, the information recorded by this operation will be - * provided to you via {@link ResourceLeakDetector}. + * Records the current access location of this object with an additional arbitrary information for + * debugging purposes. If this object is determined to be leaked, the information recorded by this + * operation will be provided to you via {@link ResourceLeakDetector}. */ @Override Payload touch(Object hint); @@ -90,4 +84,4 @@ default String getMetadataUtf8() { default String getDataUtf8() { return sliceData().toString(StandardCharsets.UTF_8); } -} \ No newline at end of file +} diff --git a/rsocket-core/src/main/java/io/rsocket/RSocketClient.java b/rsocket-core/src/main/java/io/rsocket/RSocketClient.java index 13eb75952..1e6219398 100644 --- a/rsocket-core/src/main/java/io/rsocket/RSocketClient.java +++ b/rsocket-core/src/main/java/io/rsocket/RSocketClient.java @@ -24,7 +24,6 @@ import io.rsocket.exceptions.Exceptions; import io.rsocket.internal.LimitableRequestPublisher; import io.rsocket.internal.UnboundedProcessor; - import java.nio.channels.ClosedChannelException; import java.time.Duration; import java.util.Collection; @@ -64,7 +63,8 @@ class RSocketClient implements RSocket { Function frameDecoder, Consumer errorConsumer, StreamIdSupplier streamIdSupplier) { - this(connection, frameDecoder, errorConsumer, streamIdSupplier, Duration.ZERO, Duration.ZERO, 0); + this( + connection, frameDecoder, errorConsumer, streamIdSupplier, Duration.ZERO, Duration.ZERO, 0); } RSocketClient( @@ -372,11 +372,13 @@ public Flux get() { public Frame apply(Payload payload) { final Frame requestFrame; if (firstPayload.compareAndSet(true, false)) { - requestFrame = Frame.Request.from( - streamId, requestType, payload, l); + requestFrame = + Frame.Request.from( + streamId, requestType, payload, l); } else { - requestFrame = Frame.PayloadFrame.from( - streamId, FrameType.NEXT, payload); + requestFrame = + Frame.PayloadFrame.from( + streamId, FrameType.NEXT, payload); } payload.release(); return requestFrame; diff --git a/rsocket-core/src/main/java/io/rsocket/RSocketFactory.java b/rsocket-core/src/main/java/io/rsocket/RSocketFactory.java index 2bd3d842b..33f56c349 100644 --- a/rsocket-core/src/main/java/io/rsocket/RSocketFactory.java +++ b/rsocket-core/src/main/java/io/rsocket/RSocketFactory.java @@ -28,12 +28,11 @@ import io.rsocket.transport.ClientTransport; import io.rsocket.transport.ServerTransport; import io.rsocket.util.DefaultPayload; +import io.rsocket.util.EmptyPayload; import java.time.Duration; import java.util.function.Consumer; import java.util.function.Function; import java.util.function.Supplier; - -import io.rsocket.util.EmptyPayload; import reactor.core.publisher.Mono; /** Factory for creating RSocket clients and servers. */ @@ -244,7 +243,10 @@ public Mono start() { .doOnNext( rSocket -> new RSocketServer( - multiplexer.asServerConnection(), rSocket, frameDecoder, errorConsumer)) + multiplexer.asServerConnection(), + rSocket, + frameDecoder, + errorConsumer)) .then(finalConnection.sendOne(setupFrame)) .then(wrappedRSocketClient); }); @@ -359,7 +361,8 @@ private Mono processSetupFrame( sender -> acceptor.get().accept(setupPayload, sender).map(plugins::applyServer)) .map( handler -> - new RSocketServer(multiplexer.asClientConnection(), handler, frameDecoder, errorConsumer)) + new RSocketServer( + multiplexer.asClientConnection(), handler, frameDecoder, errorConsumer)) .then(); } } diff --git a/rsocket-core/src/main/java/io/rsocket/RSocketServer.java b/rsocket-core/src/main/java/io/rsocket/RSocketServer.java index 88a5702df..7b959e756 100644 --- a/rsocket-core/src/main/java/io/rsocket/RSocketServer.java +++ b/rsocket-core/src/main/java/io/rsocket/RSocketServer.java @@ -26,7 +26,6 @@ import io.rsocket.exceptions.ApplicationException; import io.rsocket.internal.LimitableRequestPublisher; import io.rsocket.internal.UnboundedProcessor; - import java.util.Collection; import java.util.function.Consumer; import java.util.function.Function; @@ -314,7 +313,8 @@ private Mono handleRequestResponse(int streamId, Mono response) { if (payload.hasMetadata()) { flags = Frame.setFlag(flags, FLAGS_M); } - final Frame frame = Frame.PayloadFrame.from(streamId, FrameType.NEXT_COMPLETE, payload, flags); + final Frame frame = + Frame.PayloadFrame.from(streamId, FrameType.NEXT_COMPLETE, payload, flags); payload.release(); return frame; }) @@ -330,11 +330,12 @@ private Mono handleRequestResponse(int streamId, Mono response) { private Mono handleStream(int streamId, Flux response, int initialRequestN) { response - .map(payload -> { - final Frame frame = Frame.PayloadFrame.from(streamId, FrameType.NEXT, payload); - payload.release(); - return frame; - }) + .map( + payload -> { + final Frame frame = Frame.PayloadFrame.from(streamId, FrameType.NEXT, payload); + payload.release(); + return frame; + }) .transform( frameFlux -> { LimitableRequestPublisher frames = LimitableRequestPublisher.wrap(frameFlux); diff --git a/rsocket-core/src/main/java/io/rsocket/internal/SwitchTransform.java b/rsocket-core/src/main/java/io/rsocket/internal/SwitchTransform.java index 02f13b66f..6204669d5 100644 --- a/rsocket-core/src/main/java/io/rsocket/internal/SwitchTransform.java +++ b/rsocket-core/src/main/java/io/rsocket/internal/SwitchTransform.java @@ -10,71 +10,71 @@ import reactor.core.publisher.Operators; public final class SwitchTransform extends Flux { - - final Publisher source; + + final Publisher source; + final BiFunction, Publisher> transformer; + + public SwitchTransform( + Publisher source, BiFunction, Publisher> transformer) { + this.source = Objects.requireNonNull(source, "source"); + this.transformer = Objects.requireNonNull(transformer, "transformer"); + } + + @Override + public void subscribe(CoreSubscriber actual) { + Flux.from(source).subscribe(new SwitchTransformSubscriber<>(actual, transformer)); + } + + static final class SwitchTransformSubscriber implements CoreSubscriber { + @SuppressWarnings("rawtypes") + static final AtomicIntegerFieldUpdater ONCE = + AtomicIntegerFieldUpdater.newUpdater(SwitchTransformSubscriber.class, "once"); + + final CoreSubscriber actual; final BiFunction, Publisher> transformer; - - public SwitchTransform( - Publisher source, BiFunction, Publisher> transformer) { - this.source = Objects.requireNonNull(source, "source"); - this.transformer = Objects.requireNonNull(transformer, "transformer"); + final UnboundedProcessor processor = new UnboundedProcessor<>(); + Subscription s; + volatile int once; + + SwitchTransformSubscriber( + CoreSubscriber actual, + BiFunction, Publisher> transformer) { + this.actual = actual; + this.transformer = transformer; } - + @Override - public void subscribe(CoreSubscriber actual) { - Flux.from(source).subscribe(new SwitchTransformSubscriber<>(actual, transformer)); + public void onSubscribe(Subscription s) { + if (Operators.validate(this.s, s)) { + this.s = s; + processor.onSubscribe(s); + } } - - static final class SwitchTransformSubscriber implements CoreSubscriber { - @SuppressWarnings("rawtypes") - static final AtomicIntegerFieldUpdater ONCE = - AtomicIntegerFieldUpdater.newUpdater(SwitchTransformSubscriber.class, "once"); - - final CoreSubscriber actual; - final BiFunction, Publisher> transformer; - final UnboundedProcessor processor = new UnboundedProcessor<>(); - Subscription s; - volatile int once; - - SwitchTransformSubscriber( - CoreSubscriber actual, - BiFunction, Publisher> transformer) { - this.actual = actual; - this.transformer = transformer; - } - - @Override - public void onSubscribe(Subscription s) { - if (Operators.validate(this.s, s)) { - this.s = s; - processor.onSubscribe(s); - } - } - - @Override - public void onNext(T t) { - if (once == 0 && ONCE.compareAndSet(this, 0, 1)) { - try { - Publisher result = - Objects.requireNonNull( - transformer.apply(t, processor), "The transformer returned a null value"); - Flux.from(result).subscribe(actual); - } catch (Throwable e) { - onError(Operators.onOperatorError(s, e, t, actual.currentContext())); - return; - } - } - processor.onNext(t); - } - - @Override - public void onError(Throwable t) { - processor.onError(t); - } - - @Override - public void onComplete() { - processor.onComplete(); + + @Override + public void onNext(T t) { + if (once == 0 && ONCE.compareAndSet(this, 0, 1)) { + try { + Publisher result = + Objects.requireNonNull( + transformer.apply(t, processor), "The transformer returned a null value"); + Flux.from(result).subscribe(actual); + } catch (Throwable e) { + onError(Operators.onOperatorError(s, e, t, actual.currentContext())); + return; } + } + processor.onNext(t); + } + + @Override + public void onError(Throwable t) { + processor.onError(t); + } + + @Override + public void onComplete() { + processor.onComplete(); } -} \ No newline at end of file + } +} diff --git a/rsocket-core/src/main/java/io/rsocket/lease/Lease.java b/rsocket-core/src/main/java/io/rsocket/lease/Lease.java index 24f075a95..416bc1998 100644 --- a/rsocket-core/src/main/java/io/rsocket/lease/Lease.java +++ b/rsocket-core/src/main/java/io/rsocket/lease/Lease.java @@ -16,8 +16,8 @@ package io.rsocket.lease; -import javax.annotation.Nullable; import java.nio.ByteBuffer; +import javax.annotation.Nullable; /** A contract for RSocket lease, which is sent by a request acceptor and is time bound. */ public interface Lease { diff --git a/rsocket-core/src/main/java/io/rsocket/util/ByteBufPayload.java b/rsocket-core/src/main/java/io/rsocket/util/ByteBufPayload.java index d20396852..330adfe07 100644 --- a/rsocket-core/src/main/java/io/rsocket/util/ByteBufPayload.java +++ b/rsocket-core/src/main/java/io/rsocket/util/ByteBufPayload.java @@ -23,11 +23,10 @@ import io.netty.util.AbstractReferenceCounted; import io.netty.util.Recycler; import io.rsocket.Payload; - -import javax.annotation.Nullable; import java.nio.ByteBuffer; import java.nio.CharBuffer; import java.nio.charset.Charset; +import javax.annotation.Nullable; public final class ByteBufPayload extends AbstractReferenceCounted implements Payload { private static final Recycler RECYCLER = @@ -122,19 +121,26 @@ public static Payload create(String data) { public static Payload create(String data, @Nullable String metadata) { return create( ByteBufUtil.writeUtf8(ByteBufAllocator.DEFAULT, data), - metadata == null ? null : ByteBufUtil.writeUtf8(ByteBufAllocator.DEFAULT, metadata) - ); + metadata == null ? null : ByteBufUtil.writeUtf8(ByteBufAllocator.DEFAULT, metadata)); } public static Payload create(CharSequence data, Charset dataCharset) { - return create(ByteBufUtil.encodeString(ByteBufAllocator.DEFAULT, CharBuffer.wrap(data), dataCharset), null); + return create( + ByteBufUtil.encodeString(ByteBufAllocator.DEFAULT, CharBuffer.wrap(data), dataCharset), + null); } - public static Payload create(CharSequence data, Charset dataCharset, @Nullable CharSequence metadata, Charset metadataCharset) { + public static Payload create( + CharSequence data, + Charset dataCharset, + @Nullable CharSequence metadata, + Charset metadataCharset) { return create( ByteBufUtil.encodeString(ByteBufAllocator.DEFAULT, CharBuffer.wrap(data), dataCharset), - metadata == null ? null : ByteBufUtil.encodeString(ByteBufAllocator.DEFAULT, CharBuffer.wrap(metadata), metadataCharset) - ); + metadata == null + ? null + : ByteBufUtil.encodeString( + ByteBufAllocator.DEFAULT, CharBuffer.wrap(metadata), metadataCharset)); } public static Payload create(byte[] data) { @@ -142,7 +148,8 @@ public static Payload create(byte[] data) { } public static Payload create(byte[] data, @Nullable byte[] metadata) { - return create(Unpooled.wrappedBuffer(data), metadata == null ? null : Unpooled.wrappedBuffer(metadata)); + return create( + Unpooled.wrappedBuffer(data), metadata == null ? null : Unpooled.wrappedBuffer(metadata)); } public static Payload create(ByteBuffer data) { @@ -150,7 +157,8 @@ public static Payload create(ByteBuffer data) { } public static Payload create(ByteBuffer data, @Nullable ByteBuffer metadata) { - return create(Unpooled.wrappedBuffer(data), metadata == null ? null : Unpooled.wrappedBuffer(metadata)); + return create( + Unpooled.wrappedBuffer(data), metadata == null ? null : Unpooled.wrappedBuffer(metadata)); } public static Payload create(ByteBuf data) { diff --git a/rsocket-core/src/main/java/io/rsocket/util/DefaultPayload.java b/rsocket-core/src/main/java/io/rsocket/util/DefaultPayload.java index 8962bc315..74c9f7f26 100644 --- a/rsocket-core/src/main/java/io/rsocket/util/DefaultPayload.java +++ b/rsocket-core/src/main/java/io/rsocket/util/DefaultPayload.java @@ -121,19 +121,21 @@ public static Payload create(CharSequence data) { public static Payload create(CharSequence data, @Nullable CharSequence metadata) { return create( StandardCharsets.UTF_8.encode(CharBuffer.wrap(data)), - metadata == null ? null : StandardCharsets.UTF_8.encode(CharBuffer.wrap(metadata)) - ); + metadata == null ? null : StandardCharsets.UTF_8.encode(CharBuffer.wrap(metadata))); } public static Payload create(CharSequence data, Charset dataCharset) { return create(dataCharset.encode(CharBuffer.wrap(data)), null); } - public static Payload create(CharSequence data, Charset dataCharset, @Nullable CharSequence metadata, Charset metadataCharset) { + public static Payload create( + CharSequence data, + Charset dataCharset, + @Nullable CharSequence metadata, + Charset metadataCharset) { return create( dataCharset.encode(CharBuffer.wrap(data)), - metadata == null ? null : metadataCharset.encode(CharBuffer.wrap(metadata)) - ); + metadata == null ? null : metadataCharset.encode(CharBuffer.wrap(metadata))); } public static Payload create(byte[] data) { diff --git a/rsocket-core/src/main/java/io/rsocket/util/EmptyPayload.java b/rsocket-core/src/main/java/io/rsocket/util/EmptyPayload.java index d0d8e8434..bc4528038 100644 --- a/rsocket-core/src/main/java/io/rsocket/util/EmptyPayload.java +++ b/rsocket-core/src/main/java/io/rsocket/util/EmptyPayload.java @@ -7,8 +7,7 @@ public class EmptyPayload implements Payload { public static final EmptyPayload INSTANCE = new EmptyPayload(); - private EmptyPayload() { - } + private EmptyPayload() {} @Override public boolean hasMetadata() { diff --git a/rsocket-core/src/test/java/io/rsocket/FrameTest.java b/rsocket-core/src/test/java/io/rsocket/FrameTest.java index 2518ac610..1badf6da2 100644 --- a/rsocket-core/src/test/java/io/rsocket/FrameTest.java +++ b/rsocket-core/src/test/java/io/rsocket/FrameTest.java @@ -10,7 +10,8 @@ public class FrameTest { @Test public void testFrameToString() { final Frame requestFrame = - Frame.Request.from(1, FrameType.REQUEST_RESPONSE, DefaultPayload.create("streaming in -> 0"), 1); + Frame.Request.from( + 1, FrameType.REQUEST_RESPONSE, DefaultPayload.create("streaming in -> 0"), 1); assertEquals( "Frame => Stream ID: 1 Type: REQUEST_RESPONSE Payload: data: \"streaming in -> 0\" ", requestFrame.toString()); @@ -20,7 +21,10 @@ public void testFrameToString() { public void testFrameWithMetadataToString() { final Frame requestFrame = Frame.Request.from( - 1, FrameType.REQUEST_RESPONSE, DefaultPayload.create("streaming in -> 0", "metadata"), 1); + 1, + FrameType.REQUEST_RESPONSE, + DefaultPayload.create("streaming in -> 0", "metadata"), + 1); assertEquals( "Frame => Stream ID: 1 Type: REQUEST_RESPONSE Payload: metadata: \"metadata\" data: \"streaming in -> 0\" ", requestFrame.toString()); @@ -30,7 +34,10 @@ public void testFrameWithMetadataToString() { public void testPayload() { Frame frame = Frame.PayloadFrame.from( - 1, FrameType.NEXT_COMPLETE, DefaultPayload.create("Hello"), FrameHeaderFlyweight.FLAGS_C); + 1, + FrameType.NEXT_COMPLETE, + DefaultPayload.create("Hello"), + FrameHeaderFlyweight.FLAGS_C); frame.toString(); } } diff --git a/rsocket-core/src/test/java/io/rsocket/RSocketClientTest.java b/rsocket-core/src/test/java/io/rsocket/RSocketClientTest.java index 5ff0db8b6..055c62278 100644 --- a/rsocket-core/src/test/java/io/rsocket/RSocketClientTest.java +++ b/rsocket-core/src/test/java/io/rsocket/RSocketClientTest.java @@ -34,14 +34,11 @@ import io.rsocket.frame.RequestFrameFlyweight; import io.rsocket.test.util.TestSubscriber; import io.rsocket.util.DefaultPayload; - +import io.rsocket.util.EmptyPayload; import java.time.Duration; import java.util.ArrayList; import java.util.List; import java.util.stream.Collectors; - -import io.rsocket.util.EmptyPayload; -import org.junit.Ignore; import org.junit.Rule; import org.junit.Test; import org.reactivestreams.Publisher; diff --git a/rsocket-core/src/test/java/io/rsocket/RSocketServerTest.java b/rsocket-core/src/test/java/io/rsocket/RSocketServerTest.java index a0396282d..df3479321 100644 --- a/rsocket-core/src/test/java/io/rsocket/RSocketServerTest.java +++ b/rsocket-core/src/test/java/io/rsocket/RSocketServerTest.java @@ -26,12 +26,10 @@ import io.rsocket.test.util.TestDuplexConnection; import io.rsocket.test.util.TestSubscriber; import io.rsocket.util.DefaultPayload; - +import io.rsocket.util.EmptyPayload; import java.util.Collection; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.atomic.AtomicBoolean; - -import io.rsocket.util.EmptyPayload; import org.junit.Ignore; import org.junit.Rule; import org.junit.Test; @@ -130,7 +128,8 @@ public void setAcceptingSocket(RSocket acceptingSocket) { @Override protected RSocketServer newRSocket() { - return new RSocketServer(connection, acceptingSocket, DefaultPayload::create, throwable -> errors.add(throwable)); + return new RSocketServer( + connection, acceptingSocket, DefaultPayload::create, throwable -> errors.add(throwable)); } private void sendRequest(int streamId, FrameType frameType) { diff --git a/rsocket-core/src/test/java/io/rsocket/RSocketTest.java b/rsocket-core/src/test/java/io/rsocket/RSocketTest.java index 3c87cafae..49a996df8 100644 --- a/rsocket-core/src/test/java/io/rsocket/RSocketTest.java +++ b/rsocket-core/src/test/java/io/rsocket/RSocketTest.java @@ -25,10 +25,9 @@ import io.rsocket.test.util.LocalDuplexConnection; import io.rsocket.test.util.TestSubscriber; import io.rsocket.util.DefaultPayload; +import io.rsocket.util.EmptyPayload; import java.util.ArrayList; import java.util.concurrent.CountDownLatch; - -import io.rsocket.util.EmptyPayload; import org.hamcrest.MatcherAssert; import org.junit.Assert; import org.junit.Rule; @@ -77,7 +76,8 @@ public Mono requestResponse(Payload payload) { @Test(timeout = 2000) public void testChannel() throws Exception { CountDownLatch latch = new CountDownLatch(10); - Flux requests = Flux.range(0, 10).map(i -> DefaultPayload.create("streaming in -> " + i)); + Flux requests = + Flux.range(0, 10).map(i -> DefaultPayload.create("streaming in -> " + i)); Flux responses = rule.crs.requestChannel(requests); @@ -128,12 +128,15 @@ public Mono requestResponse(Payload payload) { @Override public Flux requestChannel(Publisher payloads) { Flux.from(payloads) - .map(payload -> DefaultPayload.create("server got -> [" + payload.toString() + "]")) + .map( + payload -> + DefaultPayload.create("server got -> [" + payload.toString() + "]")) .subscribe(); return Flux.range(1, 10) .map( - payload -> DefaultPayload.create("server got -> [" + payload.toString() + "]")); + payload -> + DefaultPayload.create("server got -> [" + payload.toString() + "]")); } }; diff --git a/rsocket-core/src/test/java/io/rsocket/fragmentation/FragmentationDuplexConnectionTest.java b/rsocket-core/src/test/java/io/rsocket/fragmentation/FragmentationDuplexConnectionTest.java index 1c06513de..3fbce23bc 100644 --- a/rsocket-core/src/test/java/io/rsocket/fragmentation/FragmentationDuplexConnectionTest.java +++ b/rsocket-core/src/test/java/io/rsocket/fragmentation/FragmentationDuplexConnectionTest.java @@ -26,7 +26,6 @@ import io.rsocket.Frame; import io.rsocket.FrameType; import io.rsocket.util.DefaultPayload; - import java.nio.ByteBuffer; import java.util.concurrent.ThreadLocalRandom; import org.junit.Test; @@ -119,7 +118,8 @@ public void testReassembleFragmentFrame() { ByteBuffer data = createRandomBytes(16); ByteBuffer metadata = createRandomBytes(16); Frame frame = - Frame.Request.from(1024, FrameType.REQUEST_RESPONSE, DefaultPayload.create(data, metadata), 1); + Frame.Request.from( + 1024, FrameType.REQUEST_RESPONSE, DefaultPayload.create(data, metadata), 1); FrameFragmenter frameFragmenter = new FrameFragmenter(2); Flux fragmentedFrames = frameFragmenter.fragment(frame); EmitterProcessor processor = EmitterProcessor.create(128); diff --git a/rsocket-core/src/test/java/io/rsocket/fragmentation/FrameFragmenterTest.java b/rsocket-core/src/test/java/io/rsocket/fragmentation/FrameFragmenterTest.java index 2e5a24685..5cb9b8723 100644 --- a/rsocket-core/src/test/java/io/rsocket/fragmentation/FrameFragmenterTest.java +++ b/rsocket-core/src/test/java/io/rsocket/fragmentation/FrameFragmenterTest.java @@ -19,7 +19,6 @@ import io.rsocket.Frame; import io.rsocket.FrameType; import io.rsocket.util.DefaultPayload; - import java.nio.ByteBuffer; import java.util.concurrent.ThreadLocalRandom; import org.junit.Test; diff --git a/rsocket-core/src/test/java/io/rsocket/fragmentation/FrameReassemblerTest.java b/rsocket-core/src/test/java/io/rsocket/fragmentation/FrameReassemblerTest.java index f95f8f0e6..bff6cf90e 100644 --- a/rsocket-core/src/test/java/io/rsocket/fragmentation/FrameReassemblerTest.java +++ b/rsocket-core/src/test/java/io/rsocket/fragmentation/FrameReassemblerTest.java @@ -19,7 +19,6 @@ import io.rsocket.Frame; import io.rsocket.FrameType; import io.rsocket.util.DefaultPayload; - import java.nio.ByteBuffer; import java.util.concurrent.ThreadLocalRandom; import org.junit.Test; @@ -32,7 +31,8 @@ public void testAppend() { ByteBuffer metadata = createRandomBytes(16); Frame from = - Frame.Request.from(1024, FrameType.REQUEST_RESPONSE, DefaultPayload.create(data, metadata), 1); + Frame.Request.from( + 1024, FrameType.REQUEST_RESPONSE, DefaultPayload.create(data, metadata), 1); FrameFragmenter frameFragmenter = new FrameFragmenter(2); FrameReassembler reassembler = new FrameReassembler(from); frameFragmenter.fragment(from).subscribe(reassembler::append); diff --git a/rsocket-core/src/test/java/io/rsocket/frame/RequestFrameFlyweightTest.java b/rsocket-core/src/test/java/io/rsocket/frame/RequestFrameFlyweightTest.java index 21895c469..677bb76a2 100644 --- a/rsocket-core/src/test/java/io/rsocket/frame/RequestFrameFlyweightTest.java +++ b/rsocket-core/src/test/java/io/rsocket/frame/RequestFrameFlyweightTest.java @@ -84,7 +84,8 @@ public void testEncodingWithNullMetadata() { Unpooled.copiedBuffer("d", StandardCharsets.UTF_8)); assertEquals("00000b0000000118000000000164", ByteBufUtil.hexDump(byteBuf, 0, encoded)); - Payload payload = DefaultPayload.create(Frame.from(stringToBuf("00000b0000000118000000000164"))); + Payload payload = + DefaultPayload.create(Frame.from(stringToBuf("00000b0000000118000000000164"))); assertFalse(payload.hasMetadata()); } diff --git a/rsocket-examples/src/main/java/io/rsocket/examples/transport/tcp/channel/ChannelEchoClient.java b/rsocket-examples/src/main/java/io/rsocket/examples/transport/tcp/channel/ChannelEchoClient.java index 923eadb1b..1b6077c87 100644 --- a/rsocket-examples/src/main/java/io/rsocket/examples/transport/tcp/channel/ChannelEchoClient.java +++ b/rsocket-examples/src/main/java/io/rsocket/examples/transport/tcp/channel/ChannelEchoClient.java @@ -22,7 +22,6 @@ import io.rsocket.transport.netty.client.TcpClientTransport; import io.rsocket.transport.netty.server.TcpServerTransport; import io.rsocket.util.DefaultPayload; - import java.time.Duration; import org.reactivestreams.Publisher; import reactor.core.publisher.Flux; @@ -44,7 +43,8 @@ public static void main(String[] args) { .block(); socket - .requestChannel(Flux.interval(Duration.ofMillis(1000)).map(i -> DefaultPayload.create("Hello"))) + .requestChannel( + Flux.interval(Duration.ofMillis(1000)).map(i -> DefaultPayload.create("Hello"))) .map(Payload::getDataUtf8) .doOnNext(System.out::println) .take(10) diff --git a/rsocket-examples/src/main/java/io/rsocket/examples/transport/tcp/duplex/DuplexClient.java b/rsocket-examples/src/main/java/io/rsocket/examples/transport/tcp/duplex/DuplexClient.java index 051ddc65e..d14c896a4 100644 --- a/rsocket-examples/src/main/java/io/rsocket/examples/transport/tcp/duplex/DuplexClient.java +++ b/rsocket-examples/src/main/java/io/rsocket/examples/transport/tcp/duplex/DuplexClient.java @@ -20,7 +20,6 @@ import io.rsocket.transport.netty.client.TcpClientTransport; import io.rsocket.transport.netty.server.TcpServerTransport; import io.rsocket.util.DefaultPayload; - import java.time.Duration; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; diff --git a/rsocket-examples/src/test/java/io/rsocket/integration/IntegrationTest.java b/rsocket-examples/src/test/java/io/rsocket/integration/IntegrationTest.java index 94652c810..d6d0741a6 100644 --- a/rsocket-examples/src/test/java/io/rsocket/integration/IntegrationTest.java +++ b/rsocket-examples/src/test/java/io/rsocket/integration/IntegrationTest.java @@ -16,6 +16,13 @@ package io.rsocket.integration; +import static org.hamcrest.Matchers.is; +import static org.junit.Assert.assertThat; +import static org.junit.Assert.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoMoreInteractions; + import io.rsocket.AbstractRSocket; import io.rsocket.Payload; import io.rsocket.RSocket; @@ -28,6 +35,8 @@ import io.rsocket.transport.netty.server.TcpServerTransport; import io.rsocket.util.DefaultPayload; import io.rsocket.util.RSocketProxy; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.atomic.AtomicInteger; import org.junit.After; import org.junit.Assert; import org.junit.Before; @@ -37,16 +46,6 @@ import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.atomic.AtomicInteger; - -import static org.hamcrest.Matchers.is; -import static org.junit.Assert.assertThat; -import static org.junit.Assert.assertTrue; -import static org.mockito.ArgumentMatchers.any; -import static org.mockito.Mockito.verify; -import static org.mockito.Mockito.verifyNoMoreInteractions; - public class IntegrationTest { private static final RSocketInterceptor clientPlugin; @@ -102,9 +101,10 @@ public void startup() { RSocketFactory.receive() .addServerPlugin(serverPlugin) .addConnectionPlugin(connectionPlugin) - .errorConsumer(t -> { - errorCount.incrementAndGet(); - }) + .errorConsumer( + t -> { + errorCount.incrementAndGet(); + }) .acceptor( (setup, sendingSocket) -> { sendingSocket @@ -122,7 +122,8 @@ public Mono requestResponse(Payload payload) { @Override public Flux requestStream(Payload payload) { - return Flux.range(1, 10_000).map(i -> DefaultPayload.create("data -> " + i)); + return Flux.range(1, 10_000) + .map(i -> DefaultPayload.create("data -> " + i)); } @Override diff --git a/rsocket-examples/src/test/java/io/rsocket/integration/TestingStreaming.java b/rsocket-examples/src/test/java/io/rsocket/integration/TestingStreaming.java index 31051566e..aa95272d1 100644 --- a/rsocket-examples/src/test/java/io/rsocket/integration/TestingStreaming.java +++ b/rsocket-examples/src/test/java/io/rsocket/integration/TestingStreaming.java @@ -7,7 +7,6 @@ import io.rsocket.transport.local.LocalClientTransport; import io.rsocket.transport.local.LocalServerTransport; import io.rsocket.util.DefaultPayload; - import java.time.Duration; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.Supplier; @@ -112,7 +111,8 @@ private Flux consumer(String s) { rSocket -> { AtomicInteger count = new AtomicInteger(); return Flux.range(1, 100) - .flatMap(i -> rSocket.requestStream(DefaultPayload.create("i -> " + i)).take(100), 1); + .flatMap( + i -> rSocket.requestStream(DefaultPayload.create("i -> " + i)).take(100), 1); }); } diff --git a/rsocket-load-balancer/src/test/java/io/rsocket/client/LoadBalancedRSocketMonoTest.java b/rsocket-load-balancer/src/test/java/io/rsocket/client/LoadBalancedRSocketMonoTest.java index e0c1788b4..822b70834 100644 --- a/rsocket-load-balancer/src/test/java/io/rsocket/client/LoadBalancedRSocketMonoTest.java +++ b/rsocket-load-balancer/src/test/java/io/rsocket/client/LoadBalancedRSocketMonoTest.java @@ -19,15 +19,13 @@ import io.rsocket.Payload; import io.rsocket.RSocket; import io.rsocket.client.filter.RSocketSupplier; -import io.rsocket.util.DefaultPayload; +import io.rsocket.util.EmptyPayload; import java.net.InetSocketAddress; import java.net.SocketAddress; import java.util.Arrays; import java.util.List; import java.util.concurrent.CountDownLatch; import java.util.function.Function; - -import io.rsocket.util.EmptyPayload; import org.junit.Assert; import org.junit.Test; import org.mockito.Mockito; diff --git a/rsocket-load-balancer/src/test/java/io/rsocket/client/RSocketSupplierTest.java b/rsocket-load-balancer/src/test/java/io/rsocket/client/RSocketSupplierTest.java index 40334d553..4b4e297c8 100644 --- a/rsocket-load-balancer/src/test/java/io/rsocket/client/RSocketSupplierTest.java +++ b/rsocket-load-balancer/src/test/java/io/rsocket/client/RSocketSupplierTest.java @@ -25,14 +25,11 @@ import io.rsocket.RSocket; import io.rsocket.client.filter.RSocketSupplier; import io.rsocket.test.TestSubscriber; -import io.rsocket.util.DefaultPayload; - +import io.rsocket.util.EmptyPayload; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.BiConsumer; - -import io.rsocket.util.EmptyPayload; import org.junit.Test; import org.mockito.Mockito; import org.reactivestreams.Publisher; diff --git a/rsocket-load-balancer/src/test/java/io/rsocket/client/TimeoutClientTest.java b/rsocket-load-balancer/src/test/java/io/rsocket/client/TimeoutClientTest.java index 3d73e8edc..e9afe2f5f 100644 --- a/rsocket-load-balancer/src/test/java/io/rsocket/client/TimeoutClientTest.java +++ b/rsocket-load-balancer/src/test/java/io/rsocket/client/TimeoutClientTest.java @@ -22,11 +22,8 @@ import io.rsocket.RSocket; import io.rsocket.client.filter.RSockets; import io.rsocket.exceptions.TimeoutException; -import io.rsocket.util.DefaultPayload; - -import java.time.Duration; - import io.rsocket.util.EmptyPayload; +import java.time.Duration; import org.hamcrest.MatcherAssert; import org.junit.Test; import org.reactivestreams.Subscriber; diff --git a/rsocket-test/src/main/java/io/rsocket/test/PingClient.java b/rsocket-test/src/main/java/io/rsocket/test/PingClient.java index 08ce2cc22..4d883196c 100644 --- a/rsocket-test/src/main/java/io/rsocket/test/PingClient.java +++ b/rsocket-test/src/main/java/io/rsocket/test/PingClient.java @@ -19,7 +19,6 @@ import io.rsocket.Payload; import io.rsocket.RSocket; import io.rsocket.util.ByteBufPayload; - import java.time.Duration; import org.HdrHistogram.Recorder; import reactor.core.publisher.Flux; diff --git a/rsocket-test/src/main/java/io/rsocket/test/PingHandler.java b/rsocket-test/src/main/java/io/rsocket/test/PingHandler.java index 38cdf6972..f8cc2eb08 100644 --- a/rsocket-test/src/main/java/io/rsocket/test/PingHandler.java +++ b/rsocket-test/src/main/java/io/rsocket/test/PingHandler.java @@ -22,7 +22,6 @@ import io.rsocket.RSocket; import io.rsocket.SocketAcceptor; import io.rsocket.util.ByteBufPayload; - import java.util.concurrent.ThreadLocalRandom; import reactor.core.publisher.Mono; diff --git a/rsocket-transport-aeron/src/test/java/io/rsocket/aeron/ClientServerTest.java b/rsocket-transport-aeron/src/test/java/io/rsocket/aeron/ClientServerTest.java index 289284a17..de34346a1 100644 --- a/rsocket-transport-aeron/src/test/java/io/rsocket/aeron/ClientServerTest.java +++ b/rsocket-transport-aeron/src/test/java/io/rsocket/aeron/ClientServerTest.java @@ -34,7 +34,8 @@ public class ClientServerTest { public void testFireNForget10() { long outputCount = Flux.range(1, 10) - .flatMap(i -> setup.getRSocket().fireAndForget(DefaultPayload.create("hello", "metadata"))) + .flatMap( + i -> setup.getRSocket().fireAndForget(DefaultPayload.create("hello", "metadata"))) .doOnError(Throwable::printStackTrace) .count() .block(); diff --git a/rsocket-transport-netty/src/test/java/io/rsocket/transport/netty/TcpPing.java b/rsocket-transport-netty/src/test/java/io/rsocket/transport/netty/TcpPing.java index 3f0122fd5..177e8a024 100644 --- a/rsocket-transport-netty/src/test/java/io/rsocket/transport/netty/TcpPing.java +++ b/rsocket-transport-netty/src/test/java/io/rsocket/transport/netty/TcpPing.java @@ -21,7 +21,6 @@ import io.rsocket.test.PingClient; import io.rsocket.transport.netty.client.TcpClientTransport; import java.time.Duration; - import org.HdrHistogram.Recorder; import reactor.core.publisher.Mono;