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
41 changes: 20 additions & 21 deletions pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@
<parent>
<groupId>io.scalecube</groupId>
<artifactId>scalecube-parent-pom</artifactId>
<version>0.0.17</version>
<version>0.0.19</version>
</parent>
<packaging>pom</packaging>

Expand All @@ -23,20 +23,22 @@

<properties>
<jackson.version>2.9.8</jackson.version>
<scalecube-cluster.version>2.2.3</scalecube-cluster.version>
<scalecube-cluster.version>2.2.5</scalecube-cluster.version>
<scalecube-benchmarks.version>1.2.2</scalecube-benchmarks.version>
<scalecube-config.version>0.3.4</scalecube-config.version>
<reactivestreams.version>1.0.2</reactivestreams.version>
<scalecube-config.version>0.3.6</scalecube-config.version>
<reactor.version>Californium-SR5</reactor.version>
<rsocket.version>0.11.17</rsocket.version>
<metrics.version>3.1.2</metrics.version>
<rsocket.version>0.11.16</rsocket.version>
<protostuff.version>1.6.0</protostuff.version>
<netty.version>4.1.31.Final</netty.version>
<netty.version>4.1.33.Final</netty.version>
<slf4j.version>1.7.7</slf4j.version>
<log4j.version>2.11.0</log4j.version>
<disruptor.version>3.4.2</disruptor.version>
<jsr305.version>3.0.2</jsr305.version>
<jctools.version>2.1.2</jctools.version>
<hamcrest-all.version>1.3</hamcrest-all.version>
<junit.version>5.1.1</junit.version>
<mockito.version>2.24.5</mockito.version>
<hamcrest.version>1.3</hamcrest.version>
</properties>

<modules>
Expand Down Expand Up @@ -75,13 +77,13 @@
<version>${scalecube-config.version}</version>
</dependency>



<!-- Reactive Streams -->
<!-- Reactor -->
<dependency>
<groupId>org.reactivestreams</groupId>
<artifactId>reactive-streams</artifactId>
<version>${reactivestreams.version}</version>
<groupId>io.projectreactor</groupId>
<artifactId>reactor-bom</artifactId>
<version>${reactor.version}</version>
<type>pom</type>
<scope>import</scope>
</dependency>

<!-- Logging -->
Expand Down Expand Up @@ -226,14 +228,6 @@
<artifactId>jctools-core</artifactId>
<version>${jctools.version}</version>
</dependency>

<!-- Test dependencies -->
<dependency>
<groupId>org.hamcrest</groupId>
<artifactId>hamcrest-all</artifactId>
<version>${hamcrest-all.version}</version>
<scope>test</scope>
</dependency>
</dependencies>
</dependencyManagement>

Expand All @@ -242,26 +236,31 @@
<dependency>
<groupId>org.junit.jupiter</groupId>
<artifactId>junit-jupiter-engine</artifactId>
<version>${junit.version}</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.junit.jupiter</groupId>
<artifactId>junit-jupiter-params</artifactId>
<version>${junit.version}</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.mockito</groupId>
<artifactId>mockito-junit-jupiter</artifactId>
<version>${mockito.version}</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.hamcrest</groupId>
<artifactId>hamcrest-all</artifactId>
<version>${hamcrest.version}</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.hamcrest</groupId>
<artifactId>hamcrest-core</artifactId>
<version>${hamcrest.version}</version>
<scope>test</scope>
</dependency>
<dependency>
Expand Down
11 changes: 0 additions & 11 deletions services-api/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -30,17 +30,6 @@
<groupId>org.slf4j</groupId>
<artifactId>slf4j-api</artifactId>
</dependency>

<dependency>
<groupId>org.hamcrest</groupId>
<artifactId>hamcrest-all</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>io.projectreactor</groupId>
<artifactId>reactor-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>

</project>
52 changes: 0 additions & 52 deletions services-api/src/main/java/io/scalecube/services/HeadAndTail.java

This file was deleted.

107 changes: 0 additions & 107 deletions services-api/src/test/java/io/scalecube/services/HeadAndTailTest.java

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -108,7 +108,7 @@ private Mono<Void> error(HttpServerResponse httpResponse, ServiceMessage respons
ByteBuf content =
response.hasData(ErrorData.class)
? encodeData(response.data(), response.dataFormatOrDefault())
: response.data();
: ((ByteBuf) response.data()).retain();

return httpResponse.status(status).sendObject(content).then();
}
Expand All @@ -120,7 +120,7 @@ private Mono<Void> noContent(HttpServerResponse httpResponse) {
private Mono<Void> ok(HttpServerResponse httpResponse, ServiceMessage response) {
ByteBuf content =
response.hasData(ByteBuf.class)
? response.data()
? ((ByteBuf) response.data()).retain()
: encodeData(response.data(), response.dataFormatOrDefault());

return httpResponse.status(OK).sendObject(content).then();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,6 @@
import io.rsocket.RSocket;
import io.rsocket.SocketAcceptor;
import io.rsocket.util.ByteBufPayload;
import io.scalecube.services.HeadAndTail;
import io.scalecube.services.api.ServiceMessage;
import io.scalecube.services.exceptions.BadRequestException;
import io.scalecube.services.exceptions.ServiceException;
Expand Down Expand Up @@ -72,14 +71,20 @@ public Flux<Payload> requestStream(Payload payload) {

@Override
public Flux<Payload> requestChannel(Publisher<Payload> payloads) {
return Flux.from(HeadAndTail.createFrom(Flux.from(payloads).map(this::toMessage)))
.flatMap(
pair -> {
ServiceMessage message = pair.head();
validateRequest(message);
Flux<ServiceMessage> messages = Flux.from(pair.tail()).startWith(message);
ServiceMethodInvoker methodInvoker = methodRegistry.getInvoker(message.qualifier());
return methodInvoker.invokeBidirectional(messages, ServiceMessageCodec::decodeData);
return Flux.from(payloads)
.map(this::toMessage)
.switchOnFirst(
(first, messages) -> {
if (first.hasValue()) {
ServiceMessage message = first.get();
validateRequest(message);
ServiceMethodInvoker methodInvoker =
methodRegistry.getInvoker(message.qualifier());
return methodInvoker.invokeBidirectional(
messages, ServiceMessageCodec::decodeData);
}

return messages;
})
.map(this::toPayload);
}
Expand Down
10 changes: 0 additions & 10 deletions services/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -33,16 +33,6 @@
<artifactId>jctools-core</artifactId>
</dependency>

<dependency>
<groupId>org.hamcrest</groupId>
<artifactId>hamcrest-all</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>io.projectreactor</groupId>
<artifactId>reactor-test</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>io.scalecube</groupId>
<artifactId>scalecube-services-discovery</artifactId>
Expand Down
36 changes: 20 additions & 16 deletions services/src/main/java/io/scalecube/services/ServiceCall.java
Original file line number Diff line number Diff line change
Expand Up @@ -201,23 +201,27 @@ public Flux<ServiceMessage> requestBidirectional(Publisher<ServiceMessage> publi
*/
public Flux<ServiceMessage> requestBidirectional(
Publisher<ServiceMessage> publisher, Class<?> responseType) {
return Flux.from(HeadAndTail.createFrom(publisher))
.flatMap(
pair -> {
ServiceMessage request = pair.head();
String qualifier = request.qualifier();
Flux<ServiceMessage> messages = Flux.from(pair.tail()).startWith(request);

if (methodRegistry.containsInvoker(qualifier)) { // local service.
return methodRegistry
.getInvoker(qualifier)
.invokeBidirectional(messages, ServiceMessageCodec::decodeData)
.map(this::throwIfError);
} else {
// remote service
return addressLookup(request)
.flatMapMany(address -> requestBidirectional(messages, responseType, address));
return Flux.from(publisher)
.switchOnFirst(
(first, messages) -> {
if (first.hasValue()) {
ServiceMessage request = first.get();
String qualifier = request.qualifier();

if (methodRegistry.containsInvoker(qualifier)) { // local service.
return methodRegistry
.getInvoker(qualifier)
.invokeBidirectional(messages, ServiceMessageCodec::decodeData)
.map(this::throwIfError);
} else {
// remote service
return addressLookup(request)
.flatMapMany(
address -> requestBidirectional(messages, responseType, address));
}
}

return messages;
});
}

Expand Down
Loading