Remove default blockTimeout on interface clients

See gh-30403
This commit is contained in:
Olga MaciaszekSharma
2023-04-28 17:57:44 +02:00
committed by rstoyanchev
parent e416dfdbc0
commit 033548a760
6 changed files with 33 additions and 15 deletions

View File

@@ -67,7 +67,7 @@ final class RSocketServiceMethod {
RSocketServiceMethod(
Method method, Class<?> containingClass, List<RSocketServiceArgumentResolver> argumentResolvers,
RSocketRequester rsocketRequester, @Nullable StringValueResolver embeddedValueResolver,
ReactiveAdapterRegistry reactiveRegistry, Duration blockTimeout) {
ReactiveAdapterRegistry reactiveRegistry, @Nullable Duration blockTimeout) {
this.method = method;
this.parameters = initMethodParameters(method);
@@ -125,7 +125,7 @@ final class RSocketServiceMethod {
private static Function<RSocketRequestValues, Object> initResponseFunction(
RSocketRequester requester, Method method,
ReactiveAdapterRegistry reactiveRegistry, Duration blockTimeout) {
ReactiveAdapterRegistry reactiveRegistry, @Nullable Duration blockTimeout) {
MethodParameter returnParam = new MethodParameter(method, -1);
Class<?> returnType = returnParam.getParameterType();
@@ -164,8 +164,10 @@ final class RSocketServiceMethod {
return reactiveAdapter.fromPublisher(responsePublisher);
}
return (blockForOptional ?
((Mono<?>) responsePublisher).blockOptional(blockTimeout) :
((Mono<?>) responsePublisher).block(blockTimeout));
(blockTimeout != null ? ((Mono<?>) responsePublisher).blockOptional(blockTimeout) :
((Mono<?>) responsePublisher).blockOptional()) :
(blockTimeout != null ? ((Mono<?>) responsePublisher).block(blockTimeout) :
((Mono<?>) responsePublisher).block()));
});
}

View File

@@ -59,13 +59,14 @@ public final class RSocketServiceProxyFactory {
private final ReactiveAdapterRegistry reactiveAdapterRegistry;
@Nullable
private final Duration blockTimeout;
private RSocketServiceProxyFactory(
RSocketRequester rsocketRequester, List<RSocketServiceArgumentResolver> argumentResolvers,
@Nullable StringValueResolver embeddedValueResolver,
ReactiveAdapterRegistry reactiveAdapterRegistry, Duration blockTimeout) {
ReactiveAdapterRegistry reactiveAdapterRegistry, @Nullable Duration blockTimeout) {
this.rsocketRequester = rsocketRequester;
this.argumentResolvers = argumentResolvers;
@@ -139,7 +140,7 @@ public final class RSocketServiceProxyFactory {
private ReactiveAdapterRegistry reactiveAdapterRegistry = ReactiveAdapterRegistry.getSharedInstance();
@Nullable
private Duration blockTimeout = Duration.ofSeconds(5);
private Duration blockTimeout;
private Builder() {
}
@@ -189,7 +190,8 @@ public final class RSocketServiceProxyFactory {
/**
* Configure how long to wait for a response for an HTTP service method
* with a synchronous (blocking) method signature.
* <p>By default this is 5 seconds.
* <p>By default this is {@code null},
* in which case means blocking on publishers is done without a timeout.
* @param blockTimeout the timeout value
* @return this same builder instance
*/
@@ -207,7 +209,7 @@ public final class RSocketServiceProxyFactory {
return new RSocketServiceProxyFactory(
this.rsocketRequester, initArgumentResolvers(),
this.embeddedValueResolver, this.reactiveAdapterRegistry,
(this.blockTimeout != null ? this.blockTimeout : Duration.ofSeconds(5)));
this.blockTimeout);
}
private List<RSocketServiceArgumentResolver> initArgumentResolvers() {