Skip to content

Commit 524d6d5

Browse files
NiteshKantstevegury
authored andcommitted
Introducing PublisherFunctions (#95)
* Function composition refactor * Incorporating review comments.
1 parent 9df4406 commit 524d6d5

11 files changed

Lines changed: 427 additions & 317 deletions

File tree

reactivesocket-client/src/main/java/io/reactivesocket/client/Builder.java renamed to reactivesocket-client/src/main/java/io/reactivesocket/client/ClientBuilder.java

Lines changed: 20 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,7 @@
3333
import java.util.function.Function;
3434
import java.util.stream.Collectors;
3535

36-
public class Builder {
36+
public class ClientBuilder {
3737
private static AtomicInteger counter = new AtomicInteger(0);
3838
private final String name;
3939

@@ -54,7 +54,7 @@ public class Builder {
5454

5555
private final Publisher<List<SocketAddress>> source;
5656

57-
private Builder(
57+
private ClientBuilder(
5858
String name,
5959
ScheduledExecutorService executor,
6060
long requestTimeout, TimeUnit requestTimeoutUnit,
@@ -77,8 +77,8 @@ private Builder(
7777
this.source = source;
7878
}
7979

80-
public Builder withRequestTimeout(long timeout, TimeUnit unit) {
81-
return new Builder(
80+
public ClientBuilder withRequestTimeout(long timeout, TimeUnit unit) {
81+
return new ClientBuilder(
8282
name,
8383
executor,
8484
timeout, unit,
@@ -90,8 +90,8 @@ public Builder withRequestTimeout(long timeout, TimeUnit unit) {
9090
);
9191
}
9292

93-
public Builder withConnectTimeout(long timeout, TimeUnit unit) {
94-
return new Builder(
93+
public ClientBuilder withConnectTimeout(long timeout, TimeUnit unit) {
94+
return new ClientBuilder(
9595
name,
9696
executor,
9797
requestTimeout, requestTimeoutUnit,
@@ -103,8 +103,8 @@ public Builder withConnectTimeout(long timeout, TimeUnit unit) {
103103
);
104104
}
105105

106-
public Builder withBackupRequest(double quantile) {
107-
return new Builder(
106+
public ClientBuilder withBackupRequest(double quantile) {
107+
return new ClientBuilder(
108108
name,
109109
executor,
110110
requestTimeout, requestTimeoutUnit,
@@ -116,8 +116,8 @@ public Builder withBackupRequest(double quantile) {
116116
);
117117
}
118118

119-
public Builder withExecutor(ScheduledExecutorService executor) {
120-
return new Builder(
119+
public ClientBuilder withExecutor(ScheduledExecutorService executor) {
120+
return new ClientBuilder(
121121
name,
122122
executor,
123123
requestTimeout, requestTimeoutUnit,
@@ -129,8 +129,8 @@ public Builder withExecutor(ScheduledExecutorService executor) {
129129
);
130130
}
131131

132-
public Builder withConnector(ReactiveSocketConnector<SocketAddress> connector) {
133-
return new Builder(
132+
public ClientBuilder withConnector(ReactiveSocketConnector<SocketAddress> connector) {
133+
return new ClientBuilder(
134134
name,
135135
executor,
136136
requestTimeout, requestTimeoutUnit,
@@ -142,8 +142,8 @@ public Builder withConnector(ReactiveSocketConnector<SocketAddress> connector)
142142
);
143143
}
144144

145-
public Builder withSource(Publisher<List<SocketAddress>> source) {
146-
return new Builder(
145+
public ClientBuilder withSource(Publisher<List<SocketAddress>> source) {
146+
return new ClientBuilder(
147147
name,
148148
executor,
149149
requestTimeout, requestTimeoutUnit,
@@ -155,8 +155,8 @@ public Builder withSource(Publisher<List<SocketAddress>> source) {
155155
);
156156
}
157157

158-
public Builder withRetries(int nbOfRetries, Function<Throwable, Boolean> retryThisException) {
159-
return new Builder(
158+
public ClientBuilder withRetries(int nbOfRetries, Function<Throwable, Boolean> retryThisException) {
159+
return new ClientBuilder(
160160
name,
161161
executor,
162162
requestTimeout, requestTimeoutUnit,
@@ -212,7 +212,8 @@ public void onNext(List<SocketAddress> socketAddresses) {
212212
socketAddresses.stream()
213213
.filter(sa -> !current.containsKey(sa))
214214
.map(connector::toFactory)
215-
.map(factory -> new TimeoutFactory<>(factory, connectTimeout, connectTimeoutUnit, executor))
215+
.map(factory -> factory.chain(TimeoutFactory.asChainFunction(connectTimeout, connectTimeoutUnit,
216+
executor)))
216217
.map(FailureAwareFactory::new)
217218
.forEach(factory -> current.put(factory.remote(), factory));
218219

@@ -239,8 +240,8 @@ public void onNext(List<SocketAddress> socketAddresses) {
239240
});
240241
}
241242

242-
public static Builder instance() {
243-
return new Builder(
243+
public static ClientBuilder instance() {
244+
return new ClientBuilder(
244245
"rs-loadbalancer-" + counter.incrementAndGet(),
245246
Executors.newScheduledThreadPool(4, new ThreadFactory() {
246247
@Override

reactivesocket-client/src/main/java/io/reactivesocket/client/filter/RetrySocket.java

Lines changed: 8 additions & 59 deletions
Original file line numberDiff line numberDiff line change
@@ -17,14 +17,11 @@
1717

1818
import io.reactivesocket.Payload;
1919
import io.reactivesocket.ReactiveSocket;
20+
import io.reactivesocket.internal.Publishers;
2021
import io.reactivesocket.util.ReactiveSocketProxy;
2122
import org.reactivestreams.Publisher;
22-
import org.reactivestreams.Subscriber;
23-
import org.reactivestreams.Subscription;
2423

25-
import java.util.concurrent.atomic.AtomicInteger;
2624
import java.util.function.Function;
27-
import java.util.function.Supplier;
2825

2926
public class RetrySocket extends ReactiveSocketProxy {
3027
private final int retry;
@@ -38,79 +35,31 @@ public RetrySocket(ReactiveSocket child, int retry, Function<Throwable, Boolean>
3835

3936
@Override
4037
public Publisher<Void> fireAndForget(Payload payload) {
41-
return subscriber -> child.fireAndForget(payload).subscribe(
42-
new RetrySubscriber<>(subscriber, () -> child.fireAndForget(payload))
43-
);
38+
return Publishers.retry(child.fireAndForget(payload), retry, retryThisException);
4439
}
4540

4641
@Override
4742
public Publisher<Payload> requestResponse(Payload payload) {
48-
return subscriber -> child.requestResponse(payload).subscribe(
49-
new RetrySubscriber<>(subscriber, () -> child.requestResponse(payload))
50-
);
43+
return Publishers.retry(child.requestResponse(payload), retry, retryThisException);
5144
}
5245

5346
@Override
5447
public Publisher<Payload> requestStream(Payload payload) {
55-
return subscriber -> child.requestStream(payload).subscribe(
56-
new RetrySubscriber<>(subscriber, () -> child.requestStream(payload))
57-
);
48+
return Publishers.retry(child.requestStream(payload), retry, retryThisException);
5849
}
5950

6051
@Override
6152
public Publisher<Payload> requestSubscription(Payload payload) {
62-
return subscriber -> child.requestSubscription(payload).subscribe(
63-
new RetrySubscriber<>(subscriber, () -> child.requestSubscription(payload))
64-
);
53+
return Publishers.retry(child.requestSubscription(payload), retry, retryThisException);
6554
}
6655

6756
@Override
68-
public Publisher<Payload> requestChannel(Publisher<Payload> payloads) {
69-
return subscriber -> child.requestChannel(payloads).subscribe(
70-
new RetrySubscriber<>(subscriber, () -> child.requestChannel(payloads))
71-
);
57+
public Publisher<Payload> requestChannel(Publisher<Payload> payload) {
58+
return Publishers.retry(child.requestChannel(payload), retry, retryThisException);
7259
}
7360

7461
@Override
7562
public Publisher<Void> metadataPush(Payload payload) {
76-
return subscriber -> child.metadataPush(payload).subscribe(
77-
new RetrySubscriber<>(subscriber, () -> child.metadataPush(payload))
78-
);
79-
}
80-
81-
private class RetrySubscriber<T> implements Subscriber<T> {
82-
private final Subscriber<? super T> child;
83-
private Supplier<Publisher<T>> action;
84-
private AtomicInteger budget;
85-
86-
private RetrySubscriber(Subscriber<? super T> child, Supplier<Publisher<T>> action) {
87-
this.child = child;
88-
this.action = action;
89-
this.budget = new AtomicInteger(retry);
90-
}
91-
92-
@Override
93-
public void onSubscribe(Subscription s) {
94-
child.onSubscribe(s);
95-
}
96-
97-
@Override
98-
public void onNext(T t) {
99-
child.onNext(t);
100-
}
101-
102-
@Override
103-
public void onError(Throwable t) {
104-
if (budget.decrementAndGet() > 0 && retryThisException.apply(t)) {
105-
action.get().subscribe(this);
106-
} else {
107-
child.onError(t);
108-
}
109-
}
110-
111-
@Override
112-
public void onComplete() {
113-
child.onComplete();
114-
}
63+
return Publishers.retry(child.metadataPush(payload), retry, retryThisException);
11564
}
11665
}

reactivesocket-client/src/main/java/io/reactivesocket/client/filter/TimeoutFactory.java

Lines changed: 16 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -17,52 +17,35 @@
1717

1818
import io.reactivesocket.ReactiveSocket;
1919
import io.reactivesocket.ReactiveSocketFactory;
20+
import io.reactivesocket.internal.Publishers;
21+
import io.reactivesocket.util.ReactiveSocketFactoryProxy;
2022
import org.reactivestreams.Publisher;
2123

22-
import java.util.concurrent.Executors;
2324
import java.util.concurrent.ScheduledExecutorService;
2425
import java.util.concurrent.TimeUnit;
2526
import java.util.function.Function;
2627

27-
public class TimeoutFactory<T> implements ReactiveSocketFactory<T> {
28-
private final ReactiveSocketFactory<T> child;
29-
private final ScheduledExecutorService executor;
30-
private final long timeout;
31-
private final TimeUnit unit;
28+
public class TimeoutFactory<T> extends ReactiveSocketFactoryProxy<T> {
3229

33-
public TimeoutFactory(ReactiveSocketFactory<T> child, long timeout, TimeUnit unit, ScheduledExecutorService executor) {
34-
this.child = child;
35-
this.timeout = timeout;
36-
this.unit = unit;
37-
this.executor = executor;
38-
}
30+
private final Publisher<Void> timer;
3931

40-
public TimeoutFactory(ReactiveSocketFactory<T> child, long timeout, TimeUnit unit) {
41-
this(child, timeout, unit, Executors.newScheduledThreadPool(2));
32+
public TimeoutFactory(ReactiveSocketFactory<T> child, long timeout, TimeUnit unit,
33+
ScheduledExecutorService executor) {
34+
super(child);
35+
timer = Publishers.timer(executor, timeout, unit);
4236
}
4337

4438
@Override
4539
public Publisher<ReactiveSocket> apply() {
46-
return subscriber ->
47-
child.apply().subscribe(new TimeoutSubscriber<>(subscriber, executor, timeout, unit));
48-
}
49-
50-
@Override
51-
public double availability() {
52-
return child.availability();
53-
}
54-
55-
@Override
56-
public T remote() {
57-
return child.remote();
58-
}
59-
60-
@Override
61-
public String toString() {
62-
return "TimeoutFactory(" + timeout + " " + unit.toString() + ")->" + child.toString();
40+
return Publishers.timeout(super.apply(), timer);
6341
}
6442

65-
public static <T> Function<ReactiveSocketFactory<T>, ReactiveSocketFactory<T>> filter(long timeout, TimeUnit unit) {
66-
return f -> new TimeoutFactory<>(f, timeout, unit);
43+
public static Function<Publisher<ReactiveSocket>, Publisher<ReactiveSocket>> asChainFunction(long timeout,
44+
TimeUnit unit,
45+
ScheduledExecutorService executor) {
46+
Publisher<Void> timer = Publishers.timer(executor, timeout, unit);
47+
return reactiveSocketPublisher -> {
48+
return Publishers.timeout(reactiveSocketPublisher, timer);
49+
};
6750
}
6851
}

reactivesocket-client/src/main/java/io/reactivesocket/client/filter/TimeoutSocket.java

Lines changed: 8 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -17,62 +17,43 @@
1717

1818
import io.reactivesocket.Payload;
1919
import io.reactivesocket.ReactiveSocket;
20+
import io.reactivesocket.internal.Publishers;
2021
import io.reactivesocket.util.ReactiveSocketProxy;
2122
import org.reactivestreams.Publisher;
22-
import org.reactivestreams.Subscriber;
2323

2424
import java.util.concurrent.Executors;
2525
import java.util.concurrent.ScheduledExecutorService;
2626
import java.util.concurrent.TimeUnit;
2727

2828
public class TimeoutSocket extends ReactiveSocketProxy {
29-
private final ScheduledExecutorService executor;
30-
private final ReactiveSocket child;
31-
private final long timeout;
32-
private final TimeUnit unit;
29+
private final Publisher<Void> timer;
3330

3431
public TimeoutSocket(ReactiveSocket child, long timeout, TimeUnit unit, ScheduledExecutorService executor) {
3532
super(child);
36-
this.child = child;
37-
this.timeout = timeout;
38-
this.unit = unit;
39-
this.executor = executor;
33+
timer = Publishers.timer(executor, timeout, unit);
4034
}
4135

4236
public TimeoutSocket(ReactiveSocket child, long timeout, TimeUnit unit) {
4337
this(child, timeout, unit, Executors.newScheduledThreadPool(2));
4438
}
4539

46-
@Override
47-
public Publisher<Void> fireAndForget(Payload payload) {
48-
return child.fireAndForget(payload);
49-
}
50-
5140
@Override
5241
public Publisher<Payload> requestResponse(Payload payload) {
53-
return subscriber ->
54-
child.requestResponse(payload).subscribe(wrap(subscriber));
42+
return Publishers.timeout(super.requestResponse(payload), timer);
5543
}
5644

5745
@Override
5846
public Publisher<Payload> requestStream(Payload payload) {
59-
return subscriber ->
60-
child.requestStream(payload).subscribe(wrap(subscriber));
47+
return Publishers.timeout(super.requestStream(payload), timer);
6148
}
6249

6350
@Override
6451
public Publisher<Payload> requestSubscription(Payload payload) {
65-
return subscriber ->
66-
child.requestSubscription(payload).subscribe(wrap(subscriber));
52+
return Publishers.timeout(super.requestSubscription(payload), timer);
6753
}
6854

6955
@Override
70-
public Publisher<Payload> requestChannel(Publisher<Payload> payloads) {
71-
return subscriber ->
72-
child.requestChannel(payloads).subscribe(wrap(subscriber));
73-
}
74-
75-
private <T> Subscriber<T> wrap(Subscriber<T> subscriber) {
76-
return new TimeoutSubscriber<>(subscriber, executor, timeout, unit);
56+
public Publisher<Payload> requestChannel(Publisher<Payload> payload) {
57+
return Publishers.timeout(super.requestChannel(payload), timer);
7758
}
7859
}

0 commit comments

Comments
 (0)