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() {