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 extends Void> 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 extends T> source;
+
+ final Publisher extends T> source;
+ final BiFunction, Publisher extends R>> transformer;
+
+ public SwitchTransform(
+ Publisher extends T> source, BiFunction, Publisher extends R>> transformer) {
+ this.source = Objects.requireNonNull(source, "source");
+ this.transformer = Objects.requireNonNull(transformer, "transformer");
+ }
+
+ @Override
+ public void subscribe(CoreSubscriber super R> 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 super R> actual;
final BiFunction, Publisher extends R>> transformer;
-
- public SwitchTransform(
- Publisher extends T> source, BiFunction, Publisher extends R>> 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 super R> actual,
+ BiFunction, Publisher extends R>> transformer) {
+ this.actual = actual;
+ this.transformer = transformer;
}
-
+
@Override
- public void subscribe(CoreSubscriber super R> 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 super R> actual;
- final BiFunction, Publisher extends R>> transformer;
- final UnboundedProcessor processor = new UnboundedProcessor<>();
- Subscription s;
- volatile int once;
-
- SwitchTransformSubscriber(
- CoreSubscriber super R> actual,
- BiFunction, Publisher extends R>> 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 extends R> 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 extends R> 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;