diff --git a/src/main/java/io/reactivesocket/ReactiveSocketFactory.java b/src/main/java/io/reactivesocket/ReactiveSocketFactory.java index 08a8445a7..79fb6ceec 100644 --- a/src/main/java/io/reactivesocket/ReactiveSocketFactory.java +++ b/src/main/java/io/reactivesocket/ReactiveSocketFactory.java @@ -104,16 +104,14 @@ public void onError(Throwable t) { @Override public void onComplete() { - if (!complete.get()) { - complete.set(true); + if (complete.compareAndSet(false, true)) { subscriber.onComplete(); } } }); executorService.schedule(() -> { - if (!complete.get()) { - complete.set(true); + if (complete.compareAndSet(false, true)) { subscriber.onError(new TimeoutException()); } }, timeout, timeUnit);