Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,6 @@
import org.junit.After;
import org.junit.Before;
import org.junit.Test;
import org.junit.Ignore;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

Expand Down Expand Up @@ -130,7 +129,6 @@ public void teardown() {
server.close().block();
}

@Ignore
@Test(timeout = 5_000L)
public void testRequest() {
client.requestResponse(new PayloadImpl("REQUEST", "META")).block();
Expand All @@ -140,7 +138,6 @@ public void testRequest() {
assertTrue(calledFrame);
}

@Ignore
@Test
public void testStream() throws Exception {
TestSubscriber subscriber = TestSubscriber.create();
Expand All @@ -151,7 +148,6 @@ public void testStream() throws Exception {
subscriber.assertNotComplete();
}

@Ignore
@Test(timeout = 5_000L)
public void testClose() throws ExecutionException, InterruptedException, TimeoutException {
client.close().block();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,6 @@
import org.junit.After;
import org.junit.Before;
import org.junit.Test;
import org.junit.Ignore;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.core.publisher.UnicastProcessor;
Expand Down Expand Up @@ -57,7 +56,6 @@ public void cleanup() {
server.close().block();
}

@Ignore
@Test(timeout = 5_000L)
public void testCompleteWithoutNext() throws InterruptedException {
handler =
Expand All @@ -74,7 +72,6 @@ public Flux<Payload> requestStream(Payload payload) {
assertFalse(hasElements);
}

@Ignore
@Test(timeout = 5_000L)
public void testSingleStream() throws InterruptedException {
handler =
Expand All @@ -92,7 +89,6 @@ public Flux<Payload> requestStream(Payload payload) {
assertEquals("RESPONSE", StandardCharsets.UTF_8.decode(result.getData()).toString());
}

@Ignore
@Test(timeout = 5_000L)
public void testZeroPayload() throws InterruptedException {
handler =
Expand All @@ -110,7 +106,6 @@ public Flux<Payload> requestStream(Payload payload) {
assertEquals("", StandardCharsets.UTF_8.decode(result.getData()).toString());
}

@Ignore
@Test(timeout = 5_000L)
public void testRequestResponseErrors() throws InterruptedException {
handler =
Expand Down Expand Up @@ -145,7 +140,6 @@ public Mono<Payload> requestResponse(Payload payload) {
assertEquals("SUCCESS", StandardCharsets.UTF_8.decode(response2.getData()).toString());
}

@Ignore
@Test(timeout = 5_000L)
public void testTwoConcurrentStreams() throws InterruptedException {
ConcurrentHashMap<String, UnicastProcessor<Payload>> map = new ConcurrentHashMap<>();
Expand Down
32 changes: 17 additions & 15 deletions rsocket-test/src/main/java/io/rsocket/test/ClientSetupRule.java
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import io.rsocket.transport.ServerTransport;
import io.rsocket.util.PayloadImpl;
import java.nio.charset.StandardCharsets;
import java.util.function.BiFunction;
import java.util.function.Function;
import java.util.function.Supplier;
import org.junit.rules.ExternalResource;
Expand All @@ -35,34 +36,35 @@
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

public class ClientSetupRule<T> extends ExternalResource {
public class ClientSetupRule<T, S extends Closeable> extends ExternalResource {

private Supplier<T> addressSupplier;
private Function<T, RSocket> clientConnector;
private Function<T, Closeable> serverInit;
private BiFunction<T, S, RSocket> clientConnector;
private Function<T, S> serverInit;

private RSocket client;

public ClientSetupRule(
Supplier<T> addressSupplier,
Function<T, ClientTransport> clientTransportSupplier,
Function<T, ServerTransport<? extends Closeable>> serverTransportSupplier) {
BiFunction<T, S, ClientTransport> clientTransportSupplier,
Function<T, ServerTransport<S>> serverTransportSupplier) {
this.addressSupplier = addressSupplier;

this.clientConnector =
address ->
RSocketFactory.connect()
.transport(clientTransportSupplier.apply(address))
.start()
.doOnError(t -> t.printStackTrace())
.block();

this.serverInit =
address ->
RSocketFactory.receive()
.acceptor((setup, sendingSocket) -> Mono.just(new TestRSocket()))
.transport(serverTransportSupplier.apply(address))
.start()
.map(s -> (S) s) // TODO fix casting
.block();

this.clientConnector =
(address, server) ->
RSocketFactory.connect()
.transport(clientTransportSupplier.apply(address, server))
.start()
.doOnError(t -> t.printStackTrace())
.block();
}

Expand All @@ -72,8 +74,8 @@ public Statement apply(Statement base, Description description) {
@Override
public void evaluate() throws Throwable {
T address = addressSupplier.get();
Closeable server = serverInit.apply(address);
client = clientConnector.apply(address);
S server = serverInit.apply(address);
client = clientConnector.apply(address, server);
base.evaluate();
server.close().block();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,12 +15,11 @@
*/
package io.rsocket.aeron.internal.reactivestreams;

import io.rsocket.Closeable;
import java.net.SocketAddress;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;
import java.util.function.Function;

import io.rsocket.Closeable;
import org.reactivestreams.Publisher;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,11 +26,8 @@
import io.rsocket.transport.ServerTransport;
import reactor.core.publisher.Mono;

import java.net.SocketAddress;
import java.util.concurrent.TimeUnit;

/** */
public class AeronServerTransport implements ServerTransport {
public class AeronServerTransport implements ServerTransport<Closeable> {
private final AeronWrapper aeronWrapper;
private final AeronSocketAddress managementSubscriptionSocket;
private final EventLoop eventLoop;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,14 +16,15 @@

package io.rsocket.aeron;

import io.rsocket.Closeable;
import io.rsocket.aeron.client.AeronClientTransport;
import io.rsocket.aeron.internal.*;
import io.rsocket.aeron.internal.reactivestreams.AeronClientChannelConnector;
import io.rsocket.aeron.internal.reactivestreams.AeronSocketAddress;
import io.rsocket.aeron.server.AeronServerTransport;
import io.rsocket.test.ClientSetupRule;

class AeronClientSetupRule extends ClientSetupRule<AeronSocketAddress> {
class AeronClientSetupRule extends ClientSetupRule<AeronSocketAddress, Closeable> {

public static final AeronSocketAddress ADDRESS =
AeronSocketAddress.create("aeron:udp", "127.0.0.1", 39790);
Expand Down Expand Up @@ -56,6 +57,6 @@ class AeronClientSetupRule extends ClientSetupRule<AeronSocketAddress> {
private static final AeronClientTransport client;

AeronClientSetupRule() {
super(() -> ADDRESS, address -> client, address -> server);
super(() -> ADDRESS, (address, server) -> client, address -> server);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -3,26 +3,25 @@
import io.rsocket.transport.ClientTransport;
import io.rsocket.transport.ServerTransport;
import io.rsocket.uri.UriHandler;

import java.net.URI;
import java.util.Optional;

public class LocalUriHandler implements UriHandler {
@Override
public Optional<ClientTransport> buildClient(URI uri) {
if (uri.getScheme().equals("local")) {
return Optional.of(LocalClientTransport.create(uri.getSchemeSpecificPart()));
}

return UriHandler.super.buildClient(uri);
@Override
public Optional<ClientTransport> buildClient(URI uri) {
if (uri.getScheme().equals("local")) {
return Optional.of(LocalClientTransport.create(uri.getSchemeSpecificPart()));
}

@Override
public Optional<ServerTransport> buildServer(URI uri) {
if (uri.getScheme().equals("local")) {
return Optional.of(LocalServerTransport.create(uri.getSchemeSpecificPart()));
}
return UriHandler.super.buildClient(uri);
}

return UriHandler.super.buildServer(uri);
@Override
public Optional<ServerTransport> buildServer(URI uri) {
if (uri.getScheme().equals("local")) {
return Optional.of(LocalServerTransport.create(uri.getSchemeSpecificPart()));
}

return UriHandler.super.buildServer(uri);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -16,35 +16,17 @@

package io.rsocket.transport.local;

import io.rsocket.RSocketFactory;
import io.rsocket.Closeable;
import io.rsocket.test.ClientSetupRule;
import io.rsocket.test.TestRSocket;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Supplier;
import reactor.core.publisher.Mono;

public class LocalClientSetupRule extends ClientSetupRule<String> {
public class LocalClientSetupRule extends ClientSetupRule<String, Closeable> {
private static final AtomicInteger uniqueNameGenerator = new AtomicInteger();

public LocalClientSetupRule() {
super(
// This needs to be called twice before it increments
// - once for the client and once for the server
new Supplier<String>() {
boolean increment = true;

@Override
public String get() {
if (increment) {
increment = false;
return "test" + uniqueNameGenerator.incrementAndGet();
} else {
increment = true;
return "test" + uniqueNameGenerator.get();
}
}
},
address -> LocalClientTransport.create(address),
() -> "test" + uniqueNameGenerator.incrementAndGet(),
(address, server) -> LocalClientTransport.create(address),
address -> LocalServerTransport.create(address));
}
}
Original file line number Diff line number Diff line change
@@ -1,12 +1,12 @@
package io.rsocket.transport.local;

import static org.junit.Assert.assertTrue;

import io.rsocket.transport.ClientTransport;
import io.rsocket.transport.ServerTransport;
import io.rsocket.uri.UriTransportRegistry;
import org.junit.Test;

import static org.junit.Assert.assertTrue;

public class LocalUriTransportRegistryTest {
@Test
public void testLocalClient() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
import io.rsocket.DuplexConnection;
import io.rsocket.transport.ClientTransport;
import io.rsocket.transport.netty.WebsocketDuplexConnection;
import java.net.InetSocketAddress;
import java.net.URI;
import reactor.core.publisher.Mono;
import reactor.ipc.netty.http.client.HttpClient;
Expand All @@ -42,6 +43,10 @@ public static WebsocketClientTransport create(String bindAddress, int port) {
return create(httpClient, "/");
}

public static WebsocketClientTransport create(InetSocketAddress address) {
return create(address.getHostName(), address.getPort());
}

public static WebsocketClientTransport create(URI uri) {
HttpClient httpClient = createClient(uri);
return create(httpClient, uri.getPath());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
import io.rsocket.transport.ServerTransport;
import io.rsocket.transport.netty.NettyDuplexConnection;
import io.rsocket.transport.netty.RSocketLengthCodec;
import java.net.InetSocketAddress;
import reactor.core.publisher.Mono;
import reactor.ipc.netty.tcp.TcpServer;

Expand All @@ -29,6 +30,11 @@ private TcpServerTransport(TcpServer server) {
this.server = server;
}

public static TcpServerTransport create(InetSocketAddress address) {
TcpServer server = TcpServer.create(address.getHostName(), address.getPort());
return create(server);
}

public static TcpServerTransport create(String bindAddress, int port) {
TcpServer server = TcpServer.create(bindAddress, port);
return create(server);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,15 +18,16 @@

import io.rsocket.test.ClientSetupRule;
import io.rsocket.transport.netty.client.TcpClientTransport;
import io.rsocket.transport.netty.server.NettyContextCloseable;
import io.rsocket.transport.netty.server.TcpServerTransport;
import java.net.InetSocketAddress;

public class TcpClientSetupRule extends ClientSetupRule<InetSocketAddress> {
public class TcpClientSetupRule extends ClientSetupRule<InetSocketAddress, NettyContextCloseable> {

public TcpClientSetupRule() {
super(
() -> InetSocketAddress.createUnresolved("localhost", 8989),
address -> TcpClientTransport.create(address.getHostName(), address.getPort()),
address -> TcpServerTransport.create(address.getHostName(), address.getPort()));
() -> InetSocketAddress.createUnresolved("localhost", 0),
(address, server) -> TcpClientTransport.create(server.address()),
address -> TcpServerTransport.create(address));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -18,15 +18,17 @@

import io.rsocket.test.ClientSetupRule;
import io.rsocket.transport.netty.client.WebsocketClientTransport;
import io.rsocket.transport.netty.server.NettyContextCloseable;
import io.rsocket.transport.netty.server.WebsocketServerTransport;
import java.net.InetSocketAddress;

public class WebsocketClientSetupRule extends ClientSetupRule<InetSocketAddress> {
public class WebsocketClientSetupRule
extends ClientSetupRule<InetSocketAddress, NettyContextCloseable> {

public WebsocketClientSetupRule() {
super(
() -> InetSocketAddress.createUnresolved("localhost", 8989),
address -> WebsocketClientTransport.create(address.getHostName(), address.getPort()),
() -> InetSocketAddress.createUnresolved("localhost", 0),
(address, server) -> WebsocketClientTransport.create(server.address()),
address -> WebsocketServerTransport.create(address.getHostName(), address.getPort()));
}
}