Fix issues in RSocketMessageHandler initialization
This commit ensures getRSocketStrategies() now reflects the state of corresponding RSocketMessageHandler properties even if those change after a call to setRSocketStrategies. RSocketMessageHandler has default Encoder/Decoder initializations consistent with the recent changes to RSocketStrategies.
This commit is contained in:
@@ -36,7 +36,11 @@ import org.springframework.context.ConfigurableApplicationContext;
|
||||
import org.springframework.context.EmbeddedValueResolverAware;
|
||||
import org.springframework.core.KotlinDetector;
|
||||
import org.springframework.core.annotation.AnnotatedElementUtils;
|
||||
import org.springframework.core.codec.ByteArrayDecoder;
|
||||
import org.springframework.core.codec.ByteBufferDecoder;
|
||||
import org.springframework.core.codec.DataBufferDecoder;
|
||||
import org.springframework.core.codec.Decoder;
|
||||
import org.springframework.core.codec.StringDecoder;
|
||||
import org.springframework.core.convert.ConversionService;
|
||||
import org.springframework.format.support.DefaultFormattingConversionService;
|
||||
import org.springframework.lang.Nullable;
|
||||
@@ -99,6 +103,10 @@ public class MessageMappingMessageHandler extends AbstractMethodMessageHandler<C
|
||||
|
||||
public MessageMappingMessageHandler() {
|
||||
setHandlerPredicate(type -> AnnotatedElementUtils.hasAnnotation(type, Controller.class));
|
||||
this.decoders.add(StringDecoder.allMimeTypes());
|
||||
this.decoders.add(new ByteBufferDecoder());
|
||||
this.decoders.add(new ByteArrayDecoder());
|
||||
this.decoders.add(new DataBufferDecoder());
|
||||
}
|
||||
|
||||
|
||||
@@ -106,6 +114,7 @@ public class MessageMappingMessageHandler extends AbstractMethodMessageHandler<C
|
||||
* Configure the decoders to use for incoming payloads.
|
||||
*/
|
||||
public void setDecoders(List<? extends Decoder<?>> decoders) {
|
||||
this.decoders.clear();
|
||||
this.decoders.addAll(decoders);
|
||||
}
|
||||
|
||||
@@ -178,6 +187,19 @@ public class MessageMappingMessageHandler extends AbstractMethodMessageHandler<C
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
public void afterPropertiesSet() {
|
||||
|
||||
// Initialize RouteMatcher before parent initializes handler mappings
|
||||
if (this.routeMatcher == null) {
|
||||
AntPathMatcher pathMatcher = new AntPathMatcher();
|
||||
pathMatcher.setPathSeparator(".");
|
||||
this.routeMatcher = new SimpleRouteMatcher(pathMatcher);
|
||||
}
|
||||
|
||||
super.afterPropertiesSet();
|
||||
}
|
||||
|
||||
@Override
|
||||
protected List<? extends HandlerMethodArgumentResolver> initArgumentResolvers() {
|
||||
List<HandlerMethodArgumentResolver> resolvers = new ArrayList<>();
|
||||
@@ -203,12 +225,6 @@ public class MessageMappingMessageHandler extends AbstractMethodMessageHandler<C
|
||||
resolvers.add(new PayloadMethodArgumentResolver(
|
||||
getDecoders(), this.validator, getReactiveAdapterRegistry(), true));
|
||||
|
||||
if (this.routeMatcher == null) {
|
||||
AntPathMatcher pathMatcher = new AntPathMatcher();
|
||||
pathMatcher.setPathSeparator(".");
|
||||
this.routeMatcher = new SimpleRouteMatcher(pathMatcher);
|
||||
}
|
||||
|
||||
return resolvers;
|
||||
}
|
||||
|
||||
|
||||
@@ -39,7 +39,6 @@ import org.springframework.core.io.buffer.DataBufferFactory;
|
||||
import org.springframework.core.io.buffer.NettyDataBufferFactory;
|
||||
import org.springframework.lang.Nullable;
|
||||
import org.springframework.util.AntPathMatcher;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.MimeTypeUtils;
|
||||
import org.springframework.util.RouteMatcher;
|
||||
import org.springframework.util.SimpleRouteMatcher;
|
||||
@@ -124,6 +123,7 @@ final class DefaultRSocketStrategies implements RSocketStrategies {
|
||||
@Nullable
|
||||
private MetadataExtractor metadataExtractor;
|
||||
|
||||
@Nullable
|
||||
private ReactiveAdapterRegistry adapterRegistry = ReactiveAdapterRegistry.getSharedInstance();
|
||||
|
||||
@Nullable
|
||||
@@ -149,6 +149,8 @@ final class DefaultRSocketStrategies implements RSocketStrategies {
|
||||
DefaultRSocketStrategiesBuilder(RSocketStrategies other) {
|
||||
this.encoders.addAll(other.encoders());
|
||||
this.decoders.addAll(other.decoders());
|
||||
this.routeMatcher = other.routeMatcher();
|
||||
this.metadataExtractor = other.metadataExtractor();
|
||||
this.adapterRegistry = other.reactiveAdapterRegistry();
|
||||
this.bufferFactory = other.dataBufferFactory();
|
||||
}
|
||||
@@ -179,26 +181,25 @@ final class DefaultRSocketStrategies implements RSocketStrategies {
|
||||
}
|
||||
|
||||
@Override
|
||||
public Builder routeMatcher(RouteMatcher routeMatcher) {
|
||||
public Builder routeMatcher(@Nullable RouteMatcher routeMatcher) {
|
||||
this.routeMatcher = routeMatcher;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Builder metadataExtractor(MetadataExtractor metadataExtractor) {
|
||||
public Builder metadataExtractor(@Nullable MetadataExtractor metadataExtractor) {
|
||||
this.metadataExtractor = metadataExtractor;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Builder reactiveAdapterStrategy(ReactiveAdapterRegistry registry) {
|
||||
Assert.notNull(registry, "ReactiveAdapterRegistry is required");
|
||||
public Builder reactiveAdapterStrategy(@Nullable ReactiveAdapterRegistry registry) {
|
||||
this.adapterRegistry = registry;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Builder dataBufferFactory(DataBufferFactory bufferFactory) {
|
||||
public Builder dataBufferFactory(@Nullable DataBufferFactory bufferFactory) {
|
||||
this.bufferFactory = bufferFactory;
|
||||
return this;
|
||||
}
|
||||
@@ -210,7 +211,7 @@ final class DefaultRSocketStrategies implements RSocketStrategies {
|
||||
this.routeMatcher != null ? this.routeMatcher : initRouteMatcher(),
|
||||
this.metadataExtractor != null ? this.metadataExtractor : initMetadataExtractor(),
|
||||
this.bufferFactory != null ? this.bufferFactory : initBufferFactory(),
|
||||
this.adapterRegistry);
|
||||
this.adapterRegistry != null ? this.adapterRegistry : initReactiveAdapterRegistry());
|
||||
}
|
||||
|
||||
private RouteMatcher initRouteMatcher() {
|
||||
@@ -228,6 +229,10 @@ final class DefaultRSocketStrategies implements RSocketStrategies {
|
||||
private DataBufferFactory initBufferFactory() {
|
||||
return new NettyDataBufferFactory(PooledByteBufAllocator.DEFAULT);
|
||||
}
|
||||
|
||||
private ReactiveAdapterRegistry initReactiveAdapterRegistry() {
|
||||
return ReactiveAdapterRegistry.getSharedInstance();
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -181,7 +181,7 @@ public interface RSocketStrategies {
|
||||
* efficiency consider using the {@code PathPatternRouteMatcher} from
|
||||
* {@code spring-web} instead.
|
||||
*/
|
||||
Builder routeMatcher(RouteMatcher routeMatcher);
|
||||
Builder routeMatcher(@Nullable RouteMatcher routeMatcher);
|
||||
|
||||
/**
|
||||
* Configure a {@link MetadataExtractor} to extract the route along with
|
||||
@@ -191,7 +191,7 @@ public interface RSocketStrategies {
|
||||
* route from {@code "message/x.rsocket.routing.v0"} or
|
||||
* {@code "text/plain"} metadata entries.
|
||||
*/
|
||||
Builder metadataExtractor(MetadataExtractor metadataExtractor);
|
||||
Builder metadataExtractor(@Nullable MetadataExtractor metadataExtractor);
|
||||
|
||||
/**
|
||||
* Configure the registry for reactive type support. This can be used to
|
||||
@@ -199,7 +199,7 @@ public interface RSocketStrategies {
|
||||
* {@link org.reactivestreams.Publisher Publisher}.
|
||||
* <p>By default this {@link ReactiveAdapterRegistry#getSharedInstance()}.
|
||||
*/
|
||||
Builder reactiveAdapterStrategy(ReactiveAdapterRegistry registry);
|
||||
Builder reactiveAdapterStrategy(@Nullable ReactiveAdapterRegistry registry);
|
||||
|
||||
/**
|
||||
* Configure the DataBufferFactory to use for allocating buffers when
|
||||
@@ -216,7 +216,7 @@ public interface RSocketStrategies {
|
||||
* <p>If using {@link DefaultDataBufferFactory} instead, there is no
|
||||
* need for related config changes in RSocket.
|
||||
*/
|
||||
Builder dataBufferFactory(DataBufferFactory bufferFactory);
|
||||
Builder dataBufferFactory(@Nullable DataBufferFactory bufferFactory);
|
||||
|
||||
/**
|
||||
* Build the {@code RSocketStrategies} instance.
|
||||
|
||||
@@ -32,6 +32,10 @@ import reactor.core.publisher.Mono;
|
||||
import org.springframework.beans.BeanUtils;
|
||||
import org.springframework.core.ReactiveAdapterRegistry;
|
||||
import org.springframework.core.annotation.AnnotatedElementUtils;
|
||||
import org.springframework.core.codec.ByteArrayEncoder;
|
||||
import org.springframework.core.codec.ByteBufferEncoder;
|
||||
import org.springframework.core.codec.CharSequenceEncoder;
|
||||
import org.springframework.core.codec.DataBufferEncoder;
|
||||
import org.springframework.core.codec.Decoder;
|
||||
import org.springframework.core.codec.Encoder;
|
||||
import org.springframework.lang.Nullable;
|
||||
@@ -82,6 +86,14 @@ public class RSocketMessageHandler extends MessageMappingMessageHandler {
|
||||
private MimeType defaultMetadataMimeType = MetadataExtractor.COMPOSITE_METADATA;
|
||||
|
||||
|
||||
public RSocketMessageHandler() {
|
||||
this.encoders.add(CharSequenceEncoder.allMimeTypes());
|
||||
this.encoders.add(new ByteBufferEncoder());
|
||||
this.encoders.add(new ByteArrayEncoder());
|
||||
this.encoders.add(new DataBufferEncoder());
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* {@inheritDoc}
|
||||
* <p>If {@link #setRSocketStrategies(RSocketStrategies) rsocketStrategies}
|
||||
@@ -104,6 +116,7 @@ public class RSocketMessageHandler extends MessageMappingMessageHandler {
|
||||
* other properties.
|
||||
*/
|
||||
public void setEncoders(List<? extends Encoder<?>> encoders) {
|
||||
this.encoders.clear();
|
||||
this.encoders.addAll(encoders);
|
||||
}
|
||||
|
||||
@@ -128,22 +141,32 @@ public class RSocketMessageHandler extends MessageMappingMessageHandler {
|
||||
* </ul>
|
||||
* <p>By default if this is not set, it is initialized from the above.
|
||||
*/
|
||||
public void setRSocketStrategies(@Nullable RSocketStrategies rsocketStrategies) {
|
||||
this.rsocketStrategies = rsocketStrategies;
|
||||
if (rsocketStrategies != null) {
|
||||
setDecoders(rsocketStrategies.decoders());
|
||||
setEncoders(rsocketStrategies.encoders());
|
||||
setReactiveAdapterRegistry(rsocketStrategies.reactiveAdapterRegistry());
|
||||
}
|
||||
public void setRSocketStrategies(RSocketStrategies rsocketStrategies) {
|
||||
setDecoders(rsocketStrategies.decoders());
|
||||
setEncoders(rsocketStrategies.encoders());
|
||||
setRouteMatcher(rsocketStrategies.routeMatcher());
|
||||
setMetadataExtractor(rsocketStrategies.metadataExtractor());
|
||||
setReactiveAdapterRegistry(rsocketStrategies.reactiveAdapterRegistry());
|
||||
}
|
||||
|
||||
/**
|
||||
* Return the configured {@link RSocketStrategies}. This may be {@code null}
|
||||
* before {@link #afterPropertiesSet()} is called.
|
||||
* Return an {@link RSocketStrategies} instance initialized from the
|
||||
* corresponding properties listed under {@link #setRSocketStrategies}.
|
||||
*/
|
||||
@Nullable
|
||||
public RSocketStrategies getRSocketStrategies() {
|
||||
return this.rsocketStrategies;
|
||||
return this.rsocketStrategies != null ? this.rsocketStrategies : initRSocketStrategies();
|
||||
}
|
||||
|
||||
private RSocketStrategies initRSocketStrategies() {
|
||||
return RSocketStrategies.builder()
|
||||
.decoders(List::clear)
|
||||
.encoders(List::clear)
|
||||
.decoders(decoders -> decoders.addAll(getDecoders()))
|
||||
.encoders(encoders -> encoders.addAll(getEncoders()))
|
||||
.routeMatcher(getRouteMatcher())
|
||||
.metadataExtractor(getMetadataExtractor())
|
||||
.reactiveAdapterStrategy(getReactiveAdapterRegistry())
|
||||
.build();
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -208,7 +231,9 @@ public class RSocketMessageHandler extends MessageMappingMessageHandler {
|
||||
@Override
|
||||
public void afterPropertiesSet() {
|
||||
|
||||
// Add argument resolver before parent initializes argument resolution
|
||||
getArgumentResolverConfigurer().addCustomResolver(new RSocketRequesterMethodArgumentResolver());
|
||||
|
||||
super.afterPropertiesSet();
|
||||
|
||||
if (getMetadataExtractor() == null) {
|
||||
@@ -217,15 +242,7 @@ public class RSocketMessageHandler extends MessageMappingMessageHandler {
|
||||
setMetadataExtractor(extractor);
|
||||
}
|
||||
|
||||
if (this.rsocketStrategies == null) {
|
||||
this.rsocketStrategies = RSocketStrategies.builder()
|
||||
.decoder(getDecoders().toArray(new Decoder<?>[0]))
|
||||
.encoder(getEncoders().toArray(new Encoder<?>[0]))
|
||||
.routeMatcher(getRouteMatcher())
|
||||
.metadataExtractor(getMetadataExtractor())
|
||||
.reactiveAdapterStrategy(getReactiveAdapterRegistry())
|
||||
.build();
|
||||
}
|
||||
this.rsocketStrategies = initRSocketStrategies();
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
Reference in New Issue
Block a user