Skip to content

Commit 4267666

Browse files
authored
Merge pull request #768 from scalecube/remove-onNexDropped
Removed onNextDropped RSocketServiceTransport
2 parents 5c5c608 + 7a56512 commit 4267666

File tree

1 file changed

+0
-10
lines changed

1 file changed

+0
-10
lines changed

services-transport-parent/services-transport-rsocket/src/main/java/io/scalecube/services/transport/rsocket/RSocketServiceTransport.java

Lines changed: 0 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -6,11 +6,9 @@
66
import io.netty.channel.nio.NioEventLoopGroup;
77
import io.netty.util.concurrent.DefaultThreadFactory;
88
import io.netty.util.concurrent.Future;
9-
import io.scalecube.services.api.ServiceMessage;
109
import io.scalecube.services.transport.api.ClientTransport;
1110
import io.scalecube.services.transport.api.DataCodec;
1211
import io.scalecube.services.transport.api.HeadersCodec;
13-
import io.scalecube.services.transport.api.ReferenceCountUtil;
1412
import io.scalecube.services.transport.api.ServerTransport;
1513
import io.scalecube.services.transport.api.ServiceMessageCodec;
1614
import io.scalecube.services.transport.api.ServiceTransport;
@@ -19,7 +17,6 @@
1917
import java.util.concurrent.ThreadFactory;
2018
import java.util.function.Function;
2119
import reactor.core.publisher.Flux;
22-
import reactor.core.publisher.Hooks;
2320
import reactor.core.publisher.Mono;
2421
import reactor.netty.FutureMono;
2522
import reactor.netty.resources.LoopResources;
@@ -29,13 +26,6 @@
2926
/** RSocket service transport. */
3027
public class RSocketServiceTransport implements ServiceTransport {
3128

32-
static {
33-
Hooks.onNextDropped(
34-
obj ->
35-
ReferenceCountUtil.safestRelease(
36-
obj instanceof ServiceMessage ? ((ServiceMessage) obj).data() : obj));
37-
}
38-
3929
private int numOfWorkers = Runtime.getRuntime().availableProcessors();
4030
private HeadersCodec headersCodec;
4131
private Collection<DataCodec> dataCodecs;

0 commit comments

Comments
 (0)