From 510054c175a73414fcbfcf12e43e8054abb0a07e Mon Sep 17 00:00:00 2001 From: Mark Paluch Date: Wed, 19 May 2021 15:09:59 +0200 Subject: [PATCH] Improve CassandraTemplate creation from Session. CassandraTemplate and its reactive and asynchronous variants created with a bare Session now picks up UserTypeResolver and CodecRegistry configured at the Session level. ReactiveSession now also exposes the logged keyspace and metadata. Closes #1133 --- .../data/cassandra/ReactiveSession.java | 39 +++++++++++++++++++ .../AbstractCassandraConfiguration.java | 12 ++++-- .../core/AsyncCassandraTemplate.java | 7 +++- .../cassandra/core/CassandraTemplate.java | 7 +++- .../core/ReactiveCassandraTemplate.java | 8 +++- .../convert/MappingCassandraConverter.java | 2 +- .../DefaultBridgedReactiveSession.java | 19 +++++++++ .../core/mapping/SimpleUserTypeResolver.java | 26 +++++++++++-- .../core/AsyncCassandraTemplateUnitTests.java | 18 +++++++-- .../core/CassandraTemplateUnitTests.java | 18 +++++++-- .../ReactiveCassandraTemplateUnitTests.java | 17 ++++++-- 11 files changed, 148 insertions(+), 25 deletions(-) diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/ReactiveSession.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/ReactiveSession.java index b5776680c..082497af1 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/ReactiveSession.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/ReactiveSession.java @@ -19,12 +19,17 @@ import reactor.core.publisher.Mono; import java.io.Closeable; import java.util.Map; +import java.util.Optional; +import com.datastax.oss.driver.api.core.CqlIdentifier; +import com.datastax.oss.driver.api.core.CqlSessionBuilder; import com.datastax.oss.driver.api.core.context.DriverContext; import com.datastax.oss.driver.api.core.cql.BoundStatement; import com.datastax.oss.driver.api.core.cql.PreparedStatement; import com.datastax.oss.driver.api.core.cql.SimpleStatement; import com.datastax.oss.driver.api.core.cql.Statement; +import com.datastax.oss.driver.api.core.metadata.Metadata; +import com.datastax.oss.driver.api.core.metadata.Node; /** * A session holds connections to a Cassandra cluster, allowing it to be queried. {@link ReactiveSession} executes @@ -49,6 +54,40 @@ import com.datastax.oss.driver.api.core.cql.Statement; */ public interface ReactiveSession extends Closeable { + /** + * Returns a snapshot of the Cassandra cluster's topology and schema metadata. + *

+ * In order to provide atomic updates, this method returns an immutable object: the node list, token map, and schema + * contained in a given instance will always be consistent with each other (but note that {@link Node} itself is not + * immutable: some of its properties will be updated dynamically, in particular {@link Node#getState()}). + *

+ * As a consequence of the above, you should call this method each time you need a fresh view of the metadata. Do + * not call it once and store the result, because it is a frozen snapshot that will become stale over time. + *

+ * If a metadata refresh triggers events (such as node added/removed, or schema events), then the new version of the + * metadata is guaranteed to be visible by the time you receive these events. + *

+ * + * @return never {@code null}, but may be empty if metadata has been disabled in the configuration. + * @since 3.2.2 + */ + Metadata getMetadata(); + + /** + * The keyspace that this session is currently connected to, or {@link Optional#empty()} if this session is not + * connected to any keyspace. + *

+ * There are two ways that this can be set: before initializing the session (either with the {@code session-keyspace} + * option in the configuration, or with {@link CqlSessionBuilder#withKeyspace(CqlIdentifier)}); or at runtime, if the + * client issues a request that changes the keyspace (such as a CQL {@code USE} query). Note that this second method + * is inherently unsafe, since other requests expecting the old keyspace might be executing concurrently. Therefore it + * is highly discouraged, aside from trivial cases (such as a cqlsh-style program where requests are never + * concurrent). + * + * @since 3.2.2 + */ + Optional getKeyspace(); + /** * Whether this Session instance has been closed. *

diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/AbstractCassandraConfiguration.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/AbstractCassandraConfiguration.java index 2e37fb706..8afe3a10d 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/AbstractCassandraConfiguration.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/AbstractCassandraConfiguration.java @@ -69,13 +69,15 @@ public abstract class AbstractCassandraConfiguration extends AbstractSessionConf @Bean public CassandraConverter cassandraConverter() { + CqlSession cqlSession = getRequiredSession(); + UserTypeResolver userTypeResolver = - new SimpleUserTypeResolver(getRequiredSession(), CqlIdentifier.fromCql(getKeyspaceName())); + new SimpleUserTypeResolver(cqlSession, CqlIdentifier.fromCql(getKeyspaceName())); MappingCassandraConverter converter = new MappingCassandraConverter(requireBeanOfType(CassandraMappingContext.class)); - converter.setCodecRegistry(getRequiredSession().getContext().getCodecRegistry()); + converter.setCodecRegistry(cqlSession.getContext().getCodecRegistry()); converter.setUserTypeResolver(userTypeResolver); converter.setCustomConversions(requireBeanOfType(CassandraCustomConversions.class)); @@ -92,8 +94,10 @@ public abstract class AbstractCassandraConfiguration extends AbstractSessionConf @Bean public CassandraMappingContext cassandraMapping() throws ClassNotFoundException { + CqlSession cqlSession = getRequiredSession(); + UserTypeResolver userTypeResolver = - new SimpleUserTypeResolver(getRequiredSession(), CqlIdentifier.fromCql(getKeyspaceName())); + new SimpleUserTypeResolver(cqlSession, CqlIdentifier.fromCql(getKeyspaceName())); CassandraMappingContext mappingContext = new CassandraMappingContext(userTypeResolver, SimpleTupleTypeFactory.DEFAULT); @@ -102,7 +106,7 @@ public abstract class AbstractCassandraConfiguration extends AbstractSessionConf getBeanClassLoader().ifPresent(mappingContext::setBeanClassLoader); - mappingContext.setCodecRegistry(getRequiredSession().getContext().getCodecRegistry()); + mappingContext.setCodecRegistry(cqlSession.getContext().getCodecRegistry()); mappingContext.setCustomConversions(customConversions); mappingContext.setInitialEntitySet(getInitialEntitySet()); mappingContext.setSimpleTypeHolder(customConversions.getSimpleTypeHolder()); diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/AsyncCassandraTemplate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/AsyncCassandraTemplate.java index 90c9e46d4..587684d9f 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/AsyncCassandraTemplate.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/AsyncCassandraTemplate.java @@ -42,6 +42,7 @@ import org.springframework.data.cassandra.core.cql.session.DefaultSessionFactory import org.springframework.data.cassandra.core.cql.util.CassandraFutureAdapter; import org.springframework.data.cassandra.core.cql.util.StatementBuilder; import org.springframework.data.cassandra.core.mapping.CassandraPersistentEntity; +import org.springframework.data.cassandra.core.mapping.SimpleUserTypeResolver; import org.springframework.data.cassandra.core.mapping.event.AfterConvertEvent; import org.springframework.data.cassandra.core.mapping.event.AfterDeleteEvent; import org.springframework.data.cassandra.core.mapping.event.AfterLoadEvent; @@ -137,7 +138,7 @@ public class AsyncCassandraTemplate * @see Session */ public AsyncCassandraTemplate(CqlSession session) { - this(session, newConverter()); + this(session, newConverter(session)); } /** @@ -963,9 +964,11 @@ public class AsyncCassandraTemplate return targetType.isInterface() || targetType.isAssignableFrom(entityType) ? entityType : targetType; } - private static MappingCassandraConverter newConverter() { + private static MappingCassandraConverter newConverter(CqlSession session) { MappingCassandraConverter converter = new MappingCassandraConverter(); + converter.setUserTypeResolver(new SimpleUserTypeResolver(session)); + converter.setCodecRegistry(session.getContext().getCodecRegistry()); converter.afterPropertiesSet(); diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java index b53cbeea3..2edfa8f95 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java @@ -50,6 +50,7 @@ import org.springframework.data.cassandra.core.cql.WriteOptions; import org.springframework.data.cassandra.core.cql.session.DefaultSessionFactory; import org.springframework.data.cassandra.core.cql.util.StatementBuilder; import org.springframework.data.cassandra.core.mapping.CassandraPersistentEntity; +import org.springframework.data.cassandra.core.mapping.SimpleUserTypeResolver; import org.springframework.data.cassandra.core.mapping.event.AfterConvertEvent; import org.springframework.data.cassandra.core.mapping.event.AfterDeleteEvent; import org.springframework.data.cassandra.core.mapping.event.AfterLoadEvent; @@ -140,7 +141,7 @@ public class CassandraTemplate implements CassandraOperations, ApplicationEventP * @see Session */ public CassandraTemplate(CqlSession session) { - this(session, newConverter()); + this(session, newConverter(session)); } /** @@ -1017,9 +1018,11 @@ public class CassandraTemplate implements CassandraOperations, ApplicationEventP return targetType.isInterface() || targetType.isAssignableFrom(entityType) ? entityType : targetType; } - private static MappingCassandraConverter newConverter() { + private static MappingCassandraConverter newConverter(CqlSession session) { MappingCassandraConverter converter = new MappingCassandraConverter(); + converter.setUserTypeResolver(new SimpleUserTypeResolver(session)); + converter.setCodecRegistry(session.getContext().getCodecRegistry()); converter.afterPropertiesSet(); diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplate.java index 9bfe48028..1e33de064 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplate.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplate.java @@ -44,6 +44,7 @@ import org.springframework.data.cassandra.core.cql.*; import org.springframework.data.cassandra.core.cql.session.DefaultReactiveSessionFactory; import org.springframework.data.cassandra.core.cql.util.StatementBuilder; import org.springframework.data.cassandra.core.mapping.CassandraPersistentEntity; +import org.springframework.data.cassandra.core.mapping.SimpleUserTypeResolver; import org.springframework.data.cassandra.core.mapping.event.AfterConvertEvent; import org.springframework.data.cassandra.core.mapping.event.AfterDeleteEvent; import org.springframework.data.cassandra.core.mapping.event.AfterLoadEvent; @@ -136,7 +137,7 @@ public class ReactiveCassandraTemplate * @see Session */ public ReactiveCassandraTemplate(ReactiveSession session) { - this(session, newConverter()); + this(session, newConverter(session)); } /** @@ -969,9 +970,12 @@ public class ReactiveCassandraTemplate return targetType.isInterface() || targetType.isAssignableFrom(entityType) ? entityType : targetType; } - private static MappingCassandraConverter newConverter() { + private static MappingCassandraConverter newConverter(ReactiveSession session) { MappingCassandraConverter converter = new MappingCassandraConverter(); + converter.setUserTypeResolver(new SimpleUserTypeResolver(session::getMetadata, + session.getKeyspace().orElse(CqlIdentifier.fromCql("system")))); + converter.setCodecRegistry(session.getContext().getCodecRegistry()); converter.afterPropertiesSet(); diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/MappingCassandraConverter.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/MappingCassandraConverter.java index b01614859..369a59a32 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/MappingCassandraConverter.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/MappingCassandraConverter.java @@ -87,7 +87,7 @@ public class MappingCassandraConverter extends AbstractCassandraConverter private CodecRegistry codecRegistry; - private UserTypeResolver userTypeResolver; + private @Nullable UserTypeResolver userTypeResolver; private @Nullable ClassLoader beanClassLoader; diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultBridgedReactiveSession.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultBridgedReactiveSession.java index 34e8a1424..86494f551 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultBridgedReactiveSession.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultBridgedReactiveSession.java @@ -23,6 +23,7 @@ import reactor.core.scheduler.Scheduler; import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.concurrent.CompletionStage; import org.slf4j.Logger; @@ -32,6 +33,7 @@ import org.springframework.data.cassandra.ReactiveResultSet; import org.springframework.data.cassandra.ReactiveSession; import org.springframework.util.Assert; +import com.datastax.oss.driver.api.core.CqlIdentifier; import com.datastax.oss.driver.api.core.CqlSession; import com.datastax.oss.driver.api.core.context.DriverContext; import com.datastax.oss.driver.api.core.cql.AsyncResultSet; @@ -44,6 +46,7 @@ import com.datastax.oss.driver.api.core.cql.PreparedStatement; import com.datastax.oss.driver.api.core.cql.Row; import com.datastax.oss.driver.api.core.cql.SimpleStatement; import com.datastax.oss.driver.api.core.cql.Statement; +import com.datastax.oss.driver.api.core.metadata.Metadata; /** * Default implementation of a {@link ReactiveSession}. This implementation bridges asynchronous {@link CqlSession} @@ -87,6 +90,22 @@ public class DefaultBridgedReactiveSession implements ReactiveSession { this.session = session; } + /* (non-Javadoc) + * @see org.springframework.data.cassandra.ReactiveSession#getMetadata() + */ + @Override + public Metadata getMetadata() { + return this.session.getMetadata(); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.ReactiveSession#getKeyspace() + */ + @Override + public Optional getKeyspace() { + return this.session.getKeyspace(); + } + /* (non-Javadoc) * @see org.springframework.data.cassandra.ReactiveSession#isClosed() */ diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/mapping/SimpleUserTypeResolver.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/mapping/SimpleUserTypeResolver.java index ff95a2c84..435c19d84 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/mapping/SimpleUserTypeResolver.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/mapping/SimpleUserTypeResolver.java @@ -15,6 +15,8 @@ */ package org.springframework.data.cassandra.core.mapping; +import java.util.function.Supplier; + import org.springframework.lang.Nullable; import org.springframework.util.Assert; @@ -32,7 +34,7 @@ import com.datastax.oss.driver.api.core.type.UserDefinedType; */ public class SimpleUserTypeResolver implements UserTypeResolver { - private final CqlSession session; + private final Supplier metadataSupplier; private final CqlIdentifier keyspaceName; @@ -46,7 +48,7 @@ public class SimpleUserTypeResolver implements UserTypeResolver { Assert.notNull(session, "Session must not be null"); - this.session = session; + this.metadataSupplier = session::getMetadata; this.keyspaceName = session.getKeyspace().orElse(CqlIdentifier.fromCql("system")); } @@ -62,7 +64,23 @@ public class SimpleUserTypeResolver implements UserTypeResolver { Assert.notNull(session, "Session must not be null"); Assert.notNull(keyspaceName, "Keyspace must not be null"); - this.session = session; + this.metadataSupplier = session::getMetadata; + this.keyspaceName = keyspaceName; + } + + /** + * Create a new {@link SimpleUserTypeResolver}. + * + * @param metadataSupplier must not be {@literal null}. + * @param keyspaceName must not be {@literal null}. + * @since 3.2.2 + */ + public SimpleUserTypeResolver(Supplier metadataSupplier, CqlIdentifier keyspaceName) { + + Assert.notNull(metadataSupplier, "Metadata supplier must not be null"); + Assert.notNull(keyspaceName, "Keyspace must not be null"); + + this.metadataSupplier = metadataSupplier; this.keyspaceName = keyspaceName; } @@ -72,7 +90,7 @@ public class SimpleUserTypeResolver implements UserTypeResolver { @Nullable @Override public UserDefinedType resolveType(CqlIdentifier typeName) { - return session.getMetadata().getKeyspace(keyspaceName) // + return metadataSupplier.get().getKeyspace(keyspaceName) // .flatMap(it -> it.getUserDefinedType(typeName)) // .orElse(null); } diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplateUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplateUnitTests.java index 4f2347c69..a399eedc8 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplateUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplateUnitTests.java @@ -51,6 +51,7 @@ import org.springframework.util.concurrent.ListenableFuture; import com.datastax.oss.driver.api.core.CqlIdentifier; import com.datastax.oss.driver.api.core.CqlSession; import com.datastax.oss.driver.api.core.NoNodeAvailableException; +import com.datastax.oss.driver.api.core.context.DriverContext; import com.datastax.oss.driver.api.core.cql.AsyncResultSet; import com.datastax.oss.driver.api.core.cql.ColumnDefinition; import com.datastax.oss.driver.api.core.cql.ColumnDefinitions; @@ -59,6 +60,7 @@ import com.datastax.oss.driver.api.core.cql.SimpleStatement; import com.datastax.oss.driver.api.core.cql.Statement; import com.datastax.oss.driver.api.core.type.DataTypes; import com.datastax.oss.driver.api.core.type.codec.registry.CodecRegistry; +import com.datastax.oss.driver.internal.core.type.codec.registry.DefaultCodecRegistry; /** * Unit tests for {@link AsyncCassandraTemplate}. @@ -67,9 +69,11 @@ import com.datastax.oss.driver.api.core.type.codec.registry.CodecRegistry; */ @ExtendWith(MockitoExtension.class) @MockitoSettings(strictness = Strictness.LENIENT) -public class AsyncCassandraTemplateUnitTests { +class AsyncCassandraTemplateUnitTests { @Mock CqlSession session; + CodecRegistry codecRegistry = new DefaultCodecRegistry("foo"); + @Mock DriverContext driverContext; @Mock AsyncResultSet resultSet; @Mock Row row; @Mock ColumnDefinition columnDefinition; @@ -86,8 +90,8 @@ public class AsyncCassandraTemplateUnitTests { @BeforeEach void setUp() { - template = new AsyncCassandraTemplate(session); - template.setUsePreparedStatements(false); + when(driverContext.getCodecRegistry()).thenReturn(codecRegistry); + when(session.getContext()).thenReturn(driverContext); when(session.executeAsync(any(Statement.class))).thenReturn(new TestResultSetFuture(resultSet)); when(row.getColumnDefinitions()).thenReturn(columnDefinitions); @@ -108,9 +112,17 @@ public class AsyncCassandraTemplateUnitTests { return entity; }); + template = new AsyncCassandraTemplate(session); + template.setUsePreparedStatements(false); template.setEntityCallbacks(callbacks); } + @Test // gh-1133 + void shouldConfigureConverterFromSession() { + assertThat(template.getConverter().getCodecRegistry()).isEqualTo(session.getContext().getCodecRegistry()); + assertThat(template.getConverter()).extracting("userTypeResolver").isNotNull(); + } + @Test // DATACASS-292 void selectUsingCqlShouldReturnMappedResults() { diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateUnitTests.java index a31a3a906..fbe1df2e3 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateUnitTests.java @@ -46,6 +46,7 @@ import org.springframework.data.mapping.callback.EntityCallbacks; import com.datastax.oss.driver.api.core.CqlIdentifier; import com.datastax.oss.driver.api.core.CqlSession; import com.datastax.oss.driver.api.core.NoNodeAvailableException; +import com.datastax.oss.driver.api.core.context.DriverContext; import com.datastax.oss.driver.api.core.cql.ColumnDefinition; import com.datastax.oss.driver.api.core.cql.ColumnDefinitions; import com.datastax.oss.driver.api.core.cql.ResultSet; @@ -54,6 +55,7 @@ import com.datastax.oss.driver.api.core.cql.SimpleStatement; import com.datastax.oss.driver.api.core.cql.Statement; import com.datastax.oss.driver.api.core.type.DataTypes; import com.datastax.oss.driver.api.core.type.codec.registry.CodecRegistry; +import com.datastax.oss.driver.internal.core.type.codec.registry.DefaultCodecRegistry; /** * Unit tests for {@link CassandraTemplate}. @@ -65,11 +67,12 @@ import com.datastax.oss.driver.api.core.type.codec.registry.CodecRegistry; class CassandraTemplateUnitTests { @Mock CqlSession session; + CodecRegistry codecRegistry = new DefaultCodecRegistry("foo"); + @Mock DriverContext driverContext; @Mock ResultSet resultSet; @Mock Row row; @Mock ColumnDefinition columnDefinition; @Mock ColumnDefinitions columnDefinitions; - @Captor ArgumentCaptor statementCaptor; private CassandraTemplate template; @@ -81,9 +84,8 @@ class CassandraTemplateUnitTests { @BeforeEach void setUp() { - template = new CassandraTemplate(session); - template.setUsePreparedStatements(false); - + when(driverContext.getCodecRegistry()).thenReturn(codecRegistry); + when(session.getContext()).thenReturn(driverContext); when(session.execute(any(Statement.class))).thenReturn(resultSet); when(row.getColumnDefinitions()).thenReturn(columnDefinitions); @@ -103,9 +105,17 @@ class CassandraTemplateUnitTests { return entity; }); + template = new CassandraTemplate(session); + template.setUsePreparedStatements(false); template.setEntityCallbacks(callbacks); } + @Test // gh-1133 + void shouldConfigureConverterFromSession() { + assertThat(template.getConverter().getCodecRegistry()).isEqualTo(session.getContext().getCodecRegistry()); + assertThat(template.getConverter()).extracting("userTypeResolver").isNotNull(); + } + @Test // DATACASS-292 void selectUsingCqlShouldReturnMappedResults() { diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateUnitTests.java index aef399c9c..a306086cf 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/ReactiveCassandraTemplateUnitTests.java @@ -49,6 +49,7 @@ import org.springframework.data.mapping.callback.ReactiveEntityCallbacks; import com.datastax.oss.driver.api.core.CqlIdentifier; import com.datastax.oss.driver.api.core.NoNodeAvailableException; +import com.datastax.oss.driver.api.core.context.DriverContext; import com.datastax.oss.driver.api.core.cql.ColumnDefinition; import com.datastax.oss.driver.api.core.cql.ColumnDefinitions; import com.datastax.oss.driver.api.core.cql.Row; @@ -56,6 +57,7 @@ import com.datastax.oss.driver.api.core.cql.SimpleStatement; import com.datastax.oss.driver.api.core.cql.Statement; import com.datastax.oss.driver.api.core.type.DataTypes; import com.datastax.oss.driver.api.core.type.codec.registry.CodecRegistry; +import com.datastax.oss.driver.internal.core.type.codec.registry.DefaultCodecRegistry; /** * Unit tests for {@link ReactiveCassandraTemplate}. @@ -67,6 +69,8 @@ import com.datastax.oss.driver.api.core.type.codec.registry.CodecRegistry; class ReactiveCassandraTemplateUnitTests { @Mock ReactiveSession session; + CodecRegistry codecRegistry = new DefaultCodecRegistry("foo"); + @Mock DriverContext driverContext; @Mock ReactiveResultSet reactiveResultSet; @Mock Row row; @Mock ColumnDefinition columnDefinition; @@ -83,9 +87,8 @@ class ReactiveCassandraTemplateUnitTests { @BeforeEach void setUp() { - template = new ReactiveCassandraTemplate(session); - template.setUsePreparedStatements(false); - + when(driverContext.getCodecRegistry()).thenReturn(codecRegistry); + when(session.getContext()).thenReturn(driverContext); when(session.execute(any(Statement.class))).thenReturn(Mono.just(reactiveResultSet)); when(row.getColumnDefinitions()).thenReturn(columnDefinitions); @@ -105,9 +108,17 @@ class ReactiveCassandraTemplateUnitTests { return Mono.just(entity); }); + template = new ReactiveCassandraTemplate(session); + template.setUsePreparedStatements(false); template.setEntityCallbacks(callbacks); } + @Test // gh-1133 + void shouldConfigureConverterFromSession() { + assertThat(template.getConverter().getCodecRegistry()).isEqualTo(session.getContext().getCodecRegistry()); + assertThat(template.getConverter()).extracting("userTypeResolver").isNotNull(); + } + @Test // DATACASS-335 void selectUsingCqlShouldReturnMappedResults() {