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 df8a69603..f6873787b 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 @@ -95,7 +95,7 @@ public abstract class AbstractCassandraConfiguration extends AbstractSessionConf bean.setConverter(beanFactory.getBean(CassandraConverter.class)); bean.setSchemaAction(getSchemaAction()); bean.setKeyspacePopulator(keyspacePopulator()); - bean.setKeyspacePopulator(keyspaceCleaner()); + bean.setKeyspaceCleaner(keyspaceCleaner()); return bean; } @@ -161,7 +161,7 @@ public abstract class AbstractCassandraConfiguration extends AbstractSessionConf mappingContext.setInitialEntitySet(getInitialEntitySet()); - CustomConversions customConversions = beanFactory.getBean(CassandraCustomConversions.class); + CustomConversions customConversions = beanFactory.getBean(CustomConversions.class); mappingContext.setCustomConversions(customConversions); mappingContext.setSimpleTypeHolder(customConversions.getSimpleTypeHolder()); @@ -185,8 +185,8 @@ public abstract class AbstractCassandraConfiguration extends AbstractSessionConf /** * Return the {@link Set} of initial entity classes. Scans by default the class path using - * {@link #getEntityBasePackages()}. Can be overriden by subclasses to skip class path scanning and return a fixed set - * of entity classes. + * {@link #getEntityBasePackages()}. Can be overridden by subclasses to skip class path scanning and return a fixed + * set of entity classes. * * @return {@link Set} of initial entity classes. * @throws ClassNotFoundException if the entity scan fails. @@ -200,11 +200,9 @@ public abstract class AbstractCassandraConfiguration extends AbstractSessionConf /** * Creates a {@link CassandraAdminTemplate}. - * - * @throws Exception if the {@link com.datastax.driver.core.Session} could not be obtained. */ @Bean - public CassandraAdminTemplate cassandraTemplate() throws Exception { + public CassandraAdminTemplate cassandraTemplate() { return new CassandraAdminTemplate(getRequiredSessionFactory(), beanFactory.getBean(CassandraConverter.class)); } @@ -216,6 +214,7 @@ public abstract class AbstractCassandraConfiguration extends AbstractSessionConf @Override public void setBeanFactory(BeanFactory beanFactory) throws BeansException { this.beanFactory = beanFactory; + super.setBeanFactory(beanFactory); } /** diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/AbstractReactiveCassandraConfiguration.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/AbstractReactiveCassandraConfiguration.java index 27619a346..c20c9a0e7 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/AbstractReactiveCassandraConfiguration.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/AbstractReactiveCassandraConfiguration.java @@ -41,8 +41,8 @@ public abstract class AbstractReactiveCassandraConfiguration extends AbstractCas private @Nullable BeanFactory beanFactory; /** - * Creates a {@link ReactiveSession} object. This wraps a {@link com.datastax.driver.core.Session} to expose Cassandra - * access in a reactive style. + * Creates a {@link ReactiveSession} object. This wraps a {@link com.datastax.oss.driver.api.core.CqlSession} to + * expose Cassandra access in a reactive style. * * @return the {@link ReactiveSession}. * @see #session() @@ -93,5 +93,6 @@ public abstract class AbstractReactiveCassandraConfiguration extends AbstractCas @Override public void setBeanFactory(BeanFactory beanFactory) throws BeansException { this.beanFactory = beanFactory; + super.setBeanFactory(beanFactory); } } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/AbstractSessionConfiguration.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/AbstractSessionConfiguration.java index f74019379..7ea4fecd1 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/AbstractSessionConfiguration.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/AbstractSessionConfiguration.java @@ -141,7 +141,7 @@ public abstract class AbstractSessionConfiguration implements BeanFactoryAware { /** * Returns the list of startup scripts to be run after {@link #getKeyspaceCreations() keyspace creations} and after - * initialization. + * initialization in the {@code system} keyspace. * * @return the list of startup scripts, may be empty but never {@link null} * @deprecated since 3.0, declare a @@ -154,7 +154,7 @@ public abstract class AbstractSessionConfiguration implements BeanFactoryAware { /** * Returns the list of shutdown scripts to be run after {@link #getKeyspaceDrops() keyspace drops} and right before - * shutdown. + * shutdown in the {@code system} keyspace. * * @return the list of shutdown scripts, may be empty but never {@link null} * @deprecated since 3.0, declare a diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/CqlSessionFactoryBean.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/CqlSessionFactoryBean.java index 733a9038c..c28d9ac08 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/CqlSessionFactoryBean.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/CqlSessionFactoryBean.java @@ -33,6 +33,7 @@ import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.FactoryBean; import org.springframework.beans.factory.InitializingBean; import org.springframework.data.cassandra.core.CassandraAdminOperations; +import org.springframework.data.cassandra.core.CassandraAdminTemplate; import org.springframework.data.cassandra.core.CassandraPersistentEntitySchemaCreator; import org.springframework.data.cassandra.core.CassandraPersistentEntitySchemaDropper; import org.springframework.data.cassandra.core.convert.CassandraConverter; @@ -401,18 +402,40 @@ public class CqlSessionFactoryBean implements FactoryBean, Initializ public void afterPropertiesSet() { CqlSessionBuilder sessionBuilder = buildBuilder(); - this.systemSession = sessionBuilder.build(); + this.systemSession = buildSystemSession(sessionBuilder); initializeCluster(this.systemSession); + this.session = buildSession(sessionBuilder); + + executeScripts(getStartupScripts().stream(), this.session); + performSchemaAction(); + this.systemSession.refreshSchema(); + this.session.refreshSchema(); + } + + /** + * Build the system session. + * + * @param sessionBuilder + * @return + */ + protected CqlSession buildSystemSession(CqlSessionBuilder sessionBuilder) { + return sessionBuilder.withKeyspace("system").build(); + } + + /** + * Build the keyspace session. + * + * @param sessionBuilder + * @return + */ + protected CqlSession buildSession(CqlSessionBuilder sessionBuilder) { if (StringUtils.hasText(getKeyspaceName())) { sessionBuilder.withKeyspace(getKeyspaceName()); } - this.session = sessionBuilder.build(); - - executeScripts(getStartupScripts().stream(), this.session); - performSchemaAction(); + return sessionBuilder.build(); } /* (non-Javadoc) @@ -426,12 +449,26 @@ public class CqlSessionFactoryBean implements FactoryBean, Initializ executeScripts(getShutdownScripts().stream(), this.session); executeSpecsAndScripts(keyspaceDrops, keyspaceShutdownScripts, this.systemSession); - systemSession.close(); - session.close(); + closeSystemSession(); + closeSession(); } } - private CqlSessionBuilder buildBuilder() { + /** + * Close the regular session object. + */ + protected void closeSession() { + session.close(); + } + + /** + * Close the system session object. + */ + protected void closeSystemSession() { + systemSession.close(); + } + + protected CqlSessionBuilder buildBuilder() { Assert.hasText(this.contactPoints, "At least one server is required"); @@ -514,9 +551,9 @@ public class CqlSessionFactoryBean implements FactoryBean, Initializ * statement. */ protected void createTables(boolean drop, boolean dropUnused, boolean ifNotExists) { - // TODO - /*CassandraAdminTemplate adminTemplate = new CassandraAdminTemplate(this.session, converter); - performSchemaActions(drop, dropUnused, ifNotExists, adminTemplate);*/ + + CassandraAdminTemplate adminTemplate = new CassandraAdminTemplate(this.session, converter); + performSchemaActions(drop, dropUnused, ifNotExists, adminTemplate); } private void performSchemaActions(boolean drop, boolean dropUnused, boolean ifNotExists, diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/SessionFactoryFactoryBean.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/SessionFactoryFactoryBean.java index f9d6ccc74..3467b0a17 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/SessionFactoryFactoryBean.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/SessionFactoryFactoryBean.java @@ -131,6 +131,8 @@ public class SessionFactoryFactoryBean extends AbstractFactoryBean entity = getRequiredPersistentEntity(entityClass); StatementBuilder select = getStatementFactory() - .selectOneById(id, (source, sink) -> getConverter().write(source, sink, entity), entity.getTableName()); + .selectOneById(id, entity, entity.getTableName()); return new MappingListenableFutureAdapter<>(getAsyncCqlOperations().queryForResultSet(select.build()), - resultSet -> resultSet.remaining() > 0); + resultSet -> resultSet.one() != null); } /* (non-Javadoc) @@ -541,7 +542,7 @@ public class AsyncCassandraTemplate .select(query.limit(1), getRequiredPersistentEntity(entityClass), getTableName(entityClass)); return new MappingListenableFutureAdapter<>(getAsyncCqlOperations().queryForResultSet(select.build()), - resultSet -> resultSet.remaining() > 0); + resultSet -> resultSet.one() != null); } /* (non-Javadoc) @@ -555,8 +556,7 @@ public class AsyncCassandraTemplate CassandraPersistentEntity entity = getRequiredPersistentEntity(entityClass); CqlIdentifier tableName = entity.getTableName(); - StatementBuilder select = getStatementFactory().selectOneById(id, entity, tableName); Function mapper = getMapper(entityClass, entityClass, tableName); return new MappingListenableFutureAdapter<>( @@ -595,11 +595,16 @@ public class AsyncCassandraTemplate StatementBuilder builder = getStatementFactory().insert(entityToUse, options, persistentEntity, tableName); - return source.isVersionedEntity() ? doInsertVersioned(builder.build(), entityToUse, source, tableName) - : doInsert(builder.build(), entityToUse, source, tableName); + if (source.isVersionedEntity()) { + + builder.apply(Insert::ifNotExists); + return doInsertVersioned(builder.build(), entityToUse, source, tableName); + } + + return doInsert(builder.build(), entityToUse, source, tableName); } - private ListenableFuture> doInsertVersioned(Statement insert, T entity, + private ListenableFuture> doInsertVersioned(SimpleStatement insert, T entity, AdaptibleEntity source, CqlIdentifier tableName) { return executeSave(entity, tableName, insert, result -> { @@ -613,7 +618,8 @@ public class AsyncCassandraTemplate } @SuppressWarnings("unused") - private ListenableFuture> doInsert(Statement insert, T entity, AdaptibleEntity source, + private ListenableFuture> doInsert(SimpleStatement insert, T entity, + AdaptibleEntity source, CqlIdentifier tableName) { return executeSave(entity, tableName, insert); @@ -703,9 +709,9 @@ public class AsyncCassandraTemplate AdaptibleEntity source, CqlIdentifier tableName) { StatementBuilder delete = getStatementFactory().delete(entity, options, getConverter(), tableName); - source.appendVersionCondition(delete); + ; - return executeDelete(entity, tableName, delete.build(), result -> { + return executeDelete(entity, tableName, source.appendVersionCondition(delete).build(), result -> { if (!result.wasApplied()) { throw new OptimisticLockingFailureException( @@ -734,8 +740,7 @@ public class AsyncCassandraTemplate CassandraPersistentEntity entity = getRequiredPersistentEntity(entityClass); CqlIdentifier tableName = entity.getTableName(); - StatementBuilder builder = getStatementFactory().deleteById(id, - (source, sink) -> getConverter().write(source, sink, entity), tableName); + StatementBuilder builder = getStatementFactory().deleteById(id, entity, tableName); SimpleStatement delete = builder.build(); maybeEmitEvent(new BeforeDeleteEvent<>(delete, entityClass, tableName)); @@ -771,13 +776,13 @@ public class AsyncCassandraTemplate // ------------------------------------------------------------------------- private ListenableFuture> executeSave(T entity, CqlIdentifier tableName, - Statement statement) { + SimpleStatement statement) { return executeSave(entity, tableName, statement, ignore -> {}); } private ListenableFuture> executeSave(T entity, CqlIdentifier tableName, - Statement statement, Consumer beforeAfterSaveEvent) { + SimpleStatement statement, Consumer beforeAfterSaveEvent) { maybeEmitEvent(new BeforeSaveEvent<>(entity, tableName, statement)); T entityToSave = maybeCallBeforeSave(entity, tableName, statement); @@ -798,7 +803,7 @@ public class AsyncCassandraTemplate }); } - private ListenableFuture executeDelete(Object entity, CqlIdentifier tableName, Statement statement, + private ListenableFuture executeDelete(Object entity, CqlIdentifier tableName, SimpleStatement statement, Consumer resultConsumer) { maybeEmitEvent(new BeforeDeleteEvent<>(statement, entity.getClass(), tableName)); @@ -927,9 +932,9 @@ public class AsyncCassandraTemplate @Value class AsyncStatementCallback implements AsyncSessionCallback, CqlProvider { - @lombok.NonNull Statement statement; + @lombok.NonNull SimpleStatement statement; - AsyncStatementCallback(Statement statement) { + AsyncStatementCallback(SimpleStatement statement) { this.statement = statement; } @@ -949,7 +954,7 @@ public class AsyncCassandraTemplate */ @Override public String getCql() { - return this.statement.toString(); + return this.statement.getQuery(); } } } 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 4ba2a5b6e..be1e2ff15 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 @@ -76,6 +76,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.querybuilder.QueryBuilder; import com.datastax.oss.driver.api.querybuilder.delete.Delete; +import com.datastax.oss.driver.api.querybuilder.insert.Insert; import com.datastax.oss.driver.api.querybuilder.insert.RegularInsert; import com.datastax.oss.driver.api.querybuilder.select.Select; import com.datastax.oss.driver.api.querybuilder.truncate.Truncate; @@ -565,11 +566,9 @@ public class CassandraTemplate implements CassandraOperations, ApplicationEventP Assert.notNull(entityClass, "Entity type must not be null"); CassandraPersistentEntity entity = getRequiredPersistentEntity(entityClass); + StatementBuilder select = getStatementFactory().selectOneById(id, - (source, sink) -> getConverter().write(source, sink, entity), entity.getTableName()); - - return getCqlOperations().queryForResultSet(select.build()).iterator().hasNext(); + return getCqlOperations().queryForResultSet(select.build()).one() != null; } /* (non-Javadoc) @@ -589,7 +588,7 @@ public class CassandraTemplate implements CassandraOperations, ApplicationEventP StatementBuilder select = getStatementFactory().selectOneById(id, - (source, sink) -> getConverter().write(source, sink, entity), tableName); + StatementBuilder builder = getStatementFactory().selectOneById(id, getConverter(), - getTableName(entityClass)); + CassandraPersistentEntity entity = getRequiredPersistentEntity(entityClass); + StatementBuilder builder = getStatementFactory().selectOneById(id, getConverter(), + StatementBuilder selectOneById(Object id, EntityWriter entityWriter, + StatementBuilder builder = StatementBuilder.of(select); @@ -687,7 +686,7 @@ public class StatementFactory { .collect(Collectors.toMap(Sort.Order::getProperty, // order -> order.isAscending() ? ClusteringOrder.ASC : ClusteringOrder.DESC)); - return select.orderBy(ordering); + return statement.orderBy(ordering); }); } @@ -706,14 +705,15 @@ public class StatementFactory { return com.datastax.oss.driver.api.querybuilder.select.Selector .column(((ColumnSelector) param).getExpression()); } - return com.datastax.oss.driver.api.querybuilder.select.Selector.function(param.toString()); + return new SimpleSelector(param.toString()); }).toArray(com.datastax.oss.driver.api.querybuilder.select.Selector[]::new); return com.datastax.oss.driver.api.querybuilder.select.Selector.function(selector.getExpression(), arguments); } - return QueryBuilder.literal(selector.getExpression()); + return com.datastax.oss.driver.api.querybuilder.select.Selector + .column(CqlIdentifier.fromInternal(selector.getExpression())); } private static StatementBuilder update(CqlIdentifier table, @@ -835,14 +835,20 @@ public class StatementFactory { if (updateOp.getValue() instanceof Set) { + Collection collection = (Collection) updateOp.getValue(); + Assert.isTrue(collection.size() == 1, "RemoveOp must contain a single set element"); - return Assignment.removeSetElement(updateOp.toCqlIdentifier(), termFactory.create(updateOp.getValue())); + return Assignment.removeSetElement(updateOp.toCqlIdentifier(), + termFactory.create(collection.iterator().next())); } if (updateOp.getValue() instanceof List) { + Collection collection = (Collection) updateOp.getValue(); + Assert.isTrue(collection.size() == 1, "RemoveOp must contain a single list element"); - return Assignment.removeListElement(updateOp.toCqlIdentifier(), termFactory.create(updateOp.getValue())); + return Assignment.removeListElement(updateOp.toCqlIdentifier(), + termFactory.create(collection.iterator().next())); } return Assignment.remove(updateOp.toCqlIdentifier(), termFactory.create(updateOp.getValue())); @@ -965,7 +971,8 @@ public class StatementFactory { private static Relation toClause(CriteriaDefinition criteriaDefinition, TermFactory factory) { - CqlIdentifier columnName = criteriaDefinition.getColumnName().getCqlIdentifier().get(); + CqlIdentifier columnName = criteriaDefinition.getColumnName().getCqlIdentifier() + .orElseGet(() -> CqlIdentifier.fromInternal(criteriaDefinition.getColumnName().toCql())); Predicate predicate = criteriaDefinition.getPredicate(); @@ -1094,4 +1101,30 @@ public class StatementFactory { throw new IllegalArgumentException(String.format("Criteria %s %s %s not supported for IF Conditions", columnName, predicate.getOperator(), predicate.getValue())); } + + static class SimpleSelector implements com.datastax.oss.driver.api.querybuilder.select.Selector { + + private final String selector; + + SimpleSelector(String selector) { + this.selector = selector; + } + + @NonNull + @Override + public com.datastax.oss.driver.api.querybuilder.select.Selector as(@NonNull CqlIdentifier alias) { + throw new UnsupportedOperationException(); + } + + @Nullable + @Override + public CqlIdentifier getAlias() { + return null; + } + + @Override + public void appendTo(@NonNull StringBuilder builder) { + builder.append(selector); + } + } } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/CassandraJodaTimeConverters.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/CassandraJodaTimeConverters.java index 4cfb57e95..3d2c30e8b 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/CassandraJodaTimeConverters.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/CassandraJodaTimeConverters.java @@ -15,17 +15,20 @@ */ package org.springframework.data.cassandra.core.convert; +import java.sql.Date; import java.util.ArrayList; import java.util.Collection; import java.util.Collections; import java.util.List; +import java.util.concurrent.TimeUnit; +import org.joda.time.LocalDate; +import org.joda.time.LocalDateTime; import org.joda.time.LocalTime; import org.springframework.core.convert.converter.Converter; import org.springframework.data.cassandra.core.mapping.CassandraSimpleTypeHolder; import org.springframework.data.cassandra.core.mapping.CassandraType; -import org.springframework.data.convert.ReadingConverter; import org.springframework.data.convert.WritingConverter; import org.springframework.util.ClassUtils; @@ -60,7 +63,15 @@ public abstract class CassandraJodaTimeConverters { List> converters = new ArrayList<>(); converters.add(MillisOfDayToLocalTimeConverter.INSTANCE); - converters.add(LocalTimeToMillisOfDayConverter.INSTANCE); + + converters.add(FromJodaLocalTimeConverter.INSTANCE); + converters.add(ToJodaLocalTimeConverter.INSTANCE); + + converters.add(FromJodaLocalDateConverter.INSTANCE); + converters.add(ToJodaLocalDateConverter.INSTANCE); + + converters.add(LocalDateTimeToInstantConverter.INSTANCE); + converters.add(InstantToLocalDateTimeConverter.INSTANCE); return converters; } @@ -70,7 +81,6 @@ public abstract class CassandraJodaTimeConverters { * * @author Mark Paluch */ - @ReadingConverter public enum MillisOfDayToLocalTimeConverter implements Converter { INSTANCE; @@ -86,8 +96,6 @@ public abstract class CassandraJodaTimeConverters { * * @author Mark Paluch */ - @WritingConverter - @CassandraType(type = CassandraSimpleTypeHolder.Name.TIME) public enum LocalTimeToMillisOfDayConverter implements Converter { INSTANCE; @@ -97,4 +105,99 @@ public abstract class CassandraJodaTimeConverters { return (long) source.getMillisOfDay(); } } + + /** + * Simple singleton to convert {@link LocalTime}s to their {@link java.time.LocalTime} representation. + * + * @author Mark Paluch + */ + @WritingConverter + @CassandraType(type = CassandraSimpleTypeHolder.Name.TIME) + public enum FromJodaLocalTimeConverter implements Converter { + + INSTANCE; + + @Override + public java.time.LocalTime convert(LocalTime source) { + return java.time.LocalTime.ofNanoOfDay(TimeUnit.MILLISECONDS.toNanos(source.getMillisOfDay())); + } + } + + /** + * Simple singleton to convert {@link java.time.LocalTime}s to their {@link LocalTime} representation. + * + * @author Mark Paluch + */ + public enum ToJodaLocalTimeConverter implements Converter { + + INSTANCE; + + @Override + public LocalTime convert(java.time.LocalTime source) { + return LocalTime.fromMillisOfDay(TimeUnit.NANOSECONDS.toMillis(source.toNanoOfDay())); + } + } + + /** + * Simple singleton to convert {@link LocalTime}s to their {@link java.time.LocalDate} representation. + * + * @author Mark Paluch + */ + @WritingConverter + @CassandraType(type = CassandraSimpleTypeHolder.Name.DATE) + public enum FromJodaLocalDateConverter implements Converter { + + INSTANCE; + + @Override + public java.time.LocalDate convert(LocalDate date) { + return java.time.LocalDate.of(date.getYear(), date.getMonthOfYear(), date.getDayOfMonth()); + } + } + + /** + * Simple singleton to convert {@link java.time.LocalTime}s to their {@link LocalDate} representation. + * + * @author Mark Paluch + */ + public enum ToJodaLocalDateConverter implements Converter { + + INSTANCE; + + @Override + public LocalDate convert(java.time.LocalDate date) { + return new LocalDate(date.getYear(), date.getMonthValue(), date.getDayOfMonth()); + } + } + + /** + * Simple singleton to convert {@link LocalDateTime}s to their {@link java.time.Instant} representation. + * + * @since 3.0 + */ + @WritingConverter + public enum LocalDateTimeToInstantConverter implements Converter { + + INSTANCE; + + @Override + public java.time.Instant convert(LocalDateTime source) { + return source.toDate().toInstant(); + } + } + + /** + * Simple singleton to convert {@link java.time.LocalDateTime}s to their {@link LocalDateTime} representation. + * + * @since 3.0 + */ + public enum InstantToLocalDateTimeConverter implements Converter { + + INSTANCE; + + @Override + public LocalDateTime convert(java.time.Instant source) { + return new LocalDateTime(Date.from(source)); + } + } } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/CassandraJsr310Converters.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/CassandraJsr310Converters.java index 7ad99d8a4..f9d314668 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/CassandraJsr310Converters.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/CassandraJsr310Converters.java @@ -15,15 +15,19 @@ */ package org.springframework.data.cassandra.core.convert; +import static java.time.ZoneId.*; + +import java.time.Instant; +import java.time.LocalDate; +import java.time.LocalDateTime; import java.time.LocalTime; import java.time.temporal.ChronoField; import java.util.ArrayList; import java.util.Collection; +import java.util.Date; import java.util.List; import org.springframework.core.convert.converter.Converter; -import org.springframework.data.cassandra.core.mapping.CassandraSimpleTypeHolder; -import org.springframework.data.cassandra.core.mapping.CassandraType; import org.springframework.data.convert.ReadingConverter; import org.springframework.data.convert.WritingConverter; @@ -51,6 +55,13 @@ public abstract class CassandraJsr310Converters { converters.add(MillisOfDayToLocalTimeConverter.INSTANCE); converters.add(LocalTimeToMillisOfDayConverter.INSTANCE); + converters.add(DateToInstantConverter.INSTANCE); + converters.add(LocalDateToInstantConverter.INSTANCE); + + converters.add(LocalDateConverter.INSTANCE); + converters.add(LocalTimeConverter.INSTANCE); + converters.add(InstantConverter.INSTANCE); + return converters; } @@ -77,8 +88,7 @@ public abstract class CassandraJsr310Converters { * @author Mark Paluch * @since 2.1 */ - @WritingConverter - @CassandraType(type = CassandraSimpleTypeHolder.Name.TIME) + @ReadingConverter public enum LocalTimeToMillisOfDayConverter implements Converter { INSTANCE; @@ -88,4 +98,85 @@ public abstract class CassandraJsr310Converters { return source.getLong(ChronoField.NANO_OF_DAY); } } + + /** + * Simple singleton to convert {@link Date}s to their Cassandra {@link Instant} representation for the CQL Timestamp + * type. Used for Cassandra 3.x to 4.x driver migration where + * + * @since 3.0 + */ + @WritingConverter + public enum DateToInstantConverter implements Converter { + + INSTANCE; + + @Override + public Instant convert(Date source) { + return source.toInstant(); + } + } + + /** + * Force {@link LocalDate} to remain a {@link LocalDate}. + * + * @since 3.0 + */ + @WritingConverter + enum LocalDateConverter implements Converter { + + INSTANCE; + + @Override + public LocalDate convert(LocalDate source) { + return source; + } + } + + /** + * Force {@link LocalTime} to remain a {@link LocalTime}. + * + * @since 3.0 + */ + @WritingConverter + enum LocalTimeConverter implements Converter { + + INSTANCE; + + @Override + public LocalTime convert(LocalTime source) { + return source; + } + } + + /** + * Force {@link Instant} to remain a {@link Instant}. + * + * @since 3.0 + */ + @WritingConverter + enum InstantConverter implements Converter { + + INSTANCE; + + @Override + public Instant convert(Instant source) { + return source; + } + } + + /** + * Force {@link LocalDateTime} to remain a {@link Instant}. + * + * @since 3.0 + */ + @WritingConverter + enum LocalDateToInstantConverter implements Converter { + + INSTANCE; + + @Override + public Instant convert(LocalDateTime source) { + return source.atZone(systemDefault()).toInstant(); + } + } } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/CassandraThreeTenBackPortConverters.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/CassandraThreeTenBackPortConverters.java index fa9c25734..4c6ae51e0 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/CassandraThreeTenBackPortConverters.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/CassandraThreeTenBackPortConverters.java @@ -15,6 +15,7 @@ */ package org.springframework.data.cassandra.core.convert; +import java.time.Instant; import java.util.ArrayList; import java.util.Collection; import java.util.Collections; @@ -29,7 +30,10 @@ import org.springframework.data.convert.ThreeTenBackPortConverters; import org.springframework.data.convert.WritingConverter; import org.springframework.util.ClassUtils; +import org.threeten.bp.LocalDate; +import org.threeten.bp.LocalDateTime; import org.threeten.bp.LocalTime; +import org.threeten.bp.ZoneId; import org.threeten.bp.temporal.ChronoField; /** @@ -67,6 +71,17 @@ public abstract class CassandraThreeTenBackPortConverters { converters.add(MillisOfDayToLocalTimeConverter.INSTANCE); converters.add(LocalTimeToMillisOfDayConverter.INSTANCE); + converters.add(FromBpLocalTimeConverter.INSTANCE); + converters.add(ToBpLocalTimeConverter.INSTANCE); + + converters.add(FromBpLocalDateConverter.INSTANCE); + converters.add(ToBpLocalDateConverter.INSTANCE); + + converters.add(FromBpLocalDateTimeConverter.INSTANCE); + converters.add(ToBpLocalDateTimeConverter.INSTANCE); + + converters.add(LocalDateTimeToInstantConverter.INSTANCE); + return converters; } @@ -93,8 +108,7 @@ public abstract class CassandraThreeTenBackPortConverters { * @author Mark Paluch * @since 2.1 */ - @WritingConverter - @CassandraType(type = CassandraSimpleTypeHolder.Name.TIME) + @ReadingConverter public enum LocalTimeToMillisOfDayConverter implements Converter { INSTANCE; @@ -104,4 +118,120 @@ public abstract class CassandraThreeTenBackPortConverters { return source.getLong(ChronoField.MILLI_OF_DAY); } } + + /** + * Simple singleton to convert {@link LocalTime}s to their {@link java.time.LocalTime} representation. + * + * @since 3.0 + */ + @WritingConverter + @CassandraType(type = CassandraSimpleTypeHolder.Name.TIME) + public enum FromBpLocalTimeConverter implements Converter { + + INSTANCE; + + @Override + public java.time.LocalTime convert(LocalTime source) { + return java.time.LocalTime.ofNanoOfDay(source.toNanoOfDay()); + } + } + + /** + * Simple singleton to convert {@link java.time.LocalTime}s to their {@link LocalTime} representation. + * + * @since 3.0 + */ + @ReadingConverter + public enum ToBpLocalTimeConverter implements Converter { + + INSTANCE; + + @Override + public LocalTime convert(java.time.LocalTime source) { + return LocalTime.ofNanoOfDay(source.toNanoOfDay()); + } + } + + /** + * Simple singleton to convert {@link LocalTime}s to their {@link java.time.LocalDate} representation. + * + * @since 3.0 + */ + @WritingConverter + @CassandraType(type = CassandraSimpleTypeHolder.Name.DATE) + public enum FromBpLocalDateConverter implements Converter { + + INSTANCE; + + @Override + public java.time.LocalDate convert(LocalDate date) { + return java.time.LocalDate.of(date.getYear(), date.getMonthValue(), date.getDayOfMonth()); + } + } + + /** + * Simple singleton to convert {@link java.time.LocalTime}s to their {@link LocalDate} representation. + * + * @since 3.0 + */ + @ReadingConverter + public enum ToBpLocalDateConverter implements Converter { + + INSTANCE; + + @Override + public LocalDate convert(java.time.LocalDate date) { + return LocalDate.of(date.getYear(), date.getMonthValue(), date.getDayOfMonth()); + } + } + + /** + * Simple singleton to convert {@link LocalDateTime}s to their {@link java.time.LocalDateTime} representation. + * + * @since 3.0 + */ + @ReadingConverter + public enum FromBpLocalDateTimeConverter implements Converter { + + INSTANCE; + + @Override + public java.time.LocalDateTime convert(LocalDateTime date) { + return java.time.LocalDateTime.of(date.getYear(), date.getMonthValue(), date.getDayOfMonth(), date.getHour(), + date.getMinute(), date.getSecond(), date.getNano()); + } + } + + /** + * Simple singleton to convert {@link java.time.LocalDateTime}s to their {@link LocalDateTime} representation. + * + * @since 3.0 + */ + @ReadingConverter + public enum ToBpLocalDateTimeConverter implements Converter { + + INSTANCE; + + @Override + public LocalDateTime convert(java.time.LocalDateTime date) { + return LocalDateTime.of(date.getYear(), date.getMonthValue(), date.getDayOfMonth(), date.getHour(), + date.getMinute(), date.getSecond(), date.getNano()); + } + } + + /** + * Force {@link LocalDateTime} to remain a {@link Instant}. + * + * @since 3.0 + */ + @WritingConverter + enum LocalDateTimeToInstantConverter implements Converter { + + INSTANCE; + + @Override + public java.time.Instant convert(LocalDateTime source) { + return Instant.ofEpochMilli(source.atZone(ZoneId.systemDefault()).toInstant().toEpochMilli()); + } + } } 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 a3781140e..b9b43c3dd 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 @@ -101,6 +101,7 @@ public class MappingCassandraConverter extends AbstractCassandraConverter CassandraMappingContext mappingContext = new CassandraMappingContext(); mappingContext.setCustomConversions(new CassandraCustomConversions(Collections.emptyList())); + mappingContext.afterPropertiesSet(); return mappingContext; } @@ -593,22 +594,31 @@ public class MappingCassandraConverter extends AbstractCassandraConverter return id; } + /** + * Check custom conversions for type override or fall back to + * {@link #determineTargetType(CassandraPersistentProperty)} + * + * @param property + * @return + */ private Class getTargetType(CassandraPersistentProperty property) { + return getCustomConversions().getCustomWriteTarget(property.getType()) + .orElseGet(() -> determineTargetType(property)); + } - return getCustomConversions().getCustomWriteTarget(property.getType()).orElseGet(() -> { - - if (property.isAnnotationPresent(CassandraType.class)) { - return getPropertyTargetType(property); - } - - if (property.isCompositePrimaryKey() || property.isCollectionLike() - || getCustomConversions().isSimpleType(property.getType())) { - - return property.getType(); - } + private Class determineTargetType(CassandraPersistentProperty property) { + if (property.isAnnotationPresent(CassandraType.class)) { return getPropertyTargetType(property); - }); + } + + if (property.isCompositePrimaryKey() || property.isCollectionLike() + || getCustomConversions().isSimpleType(property.getType())) { + + return property.getType(); + } + + return getPropertyTargetType(property); } private Class getPropertyTargetType(CassandraPersistentProperty property) { @@ -638,7 +648,7 @@ public class MappingCassandraConverter extends AbstractCassandraConverter @Nullable @SuppressWarnings("unchecked") private T getWriteValue(CassandraPersistentProperty property, ConvertingPropertyAccessor propertyAccessor) { - return (T) getWriteValue(propertyAccessor.getProperty(property, (Class) getTargetType(property)), + return (T) getWriteValue(propertyAccessor.getProperty(property, (Class) determineTargetType(property)), property.getTypeInformation()); } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/RowReader.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/RowReader.java index 53dbf0822..01683a3b4 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/RowReader.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/RowReader.java @@ -149,7 +149,7 @@ class RowReader { DataType valueType = setType.getElementType(); TypeCodec typeCodec = codecRegistry.codecFor(valueType); - return row.getList(index, typeCodec.getJavaType().getRawType()); + return row.getSet(index, typeCodec.getJavaType().getRawType()); } // Map diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CassandraExceptionTranslator.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CassandraExceptionTranslator.java index 2b3aa3f09..6c6f138c3 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CassandraExceptionTranslator.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/CassandraExceptionTranslator.java @@ -87,16 +87,6 @@ public class CassandraExceptionTranslator implements CqlExceptionTranslator { return new CassandraAuthenticationException(((AuthenticationException) exception).getEndPoint(), message, exception); } - /* TODO ??? - if (exception instanceof DriverInternalError) { - return new CassandraInternalException(message, exception); - } - - if (exception instanceof InvalidTypeException) { - return new CassandraTypeMismatchException(message, exception); - } - ??? - */ if (exception instanceof ReadTimeoutException) { return new CassandraReadTimeoutException(((ReadTimeoutException) exception).wasDataPresent(), message, exception); diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/QueryOptionsUtil.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/QueryOptionsUtil.java index 2ce4cbebd..7402cf9f5 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/QueryOptionsUtil.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/QueryOptionsUtil.java @@ -15,13 +15,17 @@ */ package org.springframework.data.cassandra.core.cql; +import java.time.Duration; + import org.springframework.util.Assert; import com.datastax.oss.driver.api.core.cql.SimpleStatementBuilder; import com.datastax.oss.driver.api.core.cql.Statement; import com.datastax.oss.driver.api.querybuilder.delete.Delete; +import com.datastax.oss.driver.api.querybuilder.delete.DeleteSelection; import com.datastax.oss.driver.api.querybuilder.insert.Insert; import com.datastax.oss.driver.api.querybuilder.update.Update; +import com.datastax.oss.driver.api.querybuilder.update.UpdateStart; /** * Utility class to associate {@link QueryOptions} and {@link WriteOptions} with QueryBuilder {@link Statement}s. @@ -81,11 +85,6 @@ public abstract class QueryOptionsUtil { statementBuilder.setConsistencyLevel(queryOptions.getConsistencyLevel()); } - // TODO: - /*if (queryOptions.getRetryPolicy() != null) { - statementToUse = statementToUse.setRetryPolicy(queryOptions.getRetryPolicy()); - } */ - if (queryOptions.getPageSize() != null) { statementBuilder.setPageSize(queryOptions.getPageSize()); } @@ -141,7 +140,9 @@ public abstract class QueryOptionsUtil { Assert.notNull(delete, "Delete must not be null"); Assert.notNull(writeOptions, "WriteOptions must not be null"); - // TODO: Timestamp? TTL + if (writeOptions.getTimestamp() != null) { + delete = (Delete) ((DeleteSelection) delete).usingTimestamp(writeOptions.getTimestamp()); + } return delete; } @@ -158,8 +159,22 @@ public abstract class QueryOptionsUtil { Assert.notNull(update, "Update must not be null"); Assert.notNull(writeOptions, "WriteOptions must not be null"); - // TODO: Timestamp, TTL? + if (hasTtl(writeOptions.getTtl())) { + update = (Update) ((UpdateStart) update).usingTtl(getTtlSeconds(writeOptions.getTtl())); + } + + if (writeOptions.getTimestamp() != null) { + update = (Update) ((UpdateStart) update).usingTimestamp(writeOptions.getTimestamp()); + } return update; } + + private static int getTtlSeconds(Duration ttl) { + return Math.toIntExact(ttl.getSeconds()); + } + + private static boolean hasTtl(Duration ttl) { + return !ttl.isZero() && !ttl.isNegative(); + } } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/generator/AddColumnCqlGenerator.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/generator/AddColumnCqlGenerator.java index 98d86531e..daad75342 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/generator/AddColumnCqlGenerator.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/generator/AddColumnCqlGenerator.java @@ -37,6 +37,7 @@ public class AddColumnCqlGenerator extends ColumnChangeCqlGenerator { * @param action the builder function to be applied to the statement. * @return {@code this} {@link StatementBuilder}. */ - public StatementBuilder apply(UnaryOperator action) { + public StatementBuilder apply(Function action) { Assert.notNull(action, "BindFunction must not be null"); - queryActions.add((source, termFactory) -> action.apply(source)); + queryActions.add((source, termFactory) -> (S) action.apply(source)); return this; } @@ -168,7 +176,7 @@ public class StatementBuilder { if (parameterHandling == ParameterHandling.INLINE) { - TermFactory termFactory = value -> QueryBuilder.literal(value, codecRegistry); + TermFactory termFactory = value -> toLiteralTerms(value, codecRegistry); for (BuilderRunnable runnable : queryActions) { statement = runnable.run(statement, termFactory); @@ -214,6 +222,44 @@ public class StatementBuilder { throw new UnsupportedOperationException(String.format("ParameterHandling %s not supported", parameterHandling)); } + private static Term toLiteralTerms(@Nullable Object value, CodecRegistry codecRegistry) { + + if (value instanceof List) { + + List terms = new ArrayList<>(); + + for (Object o : (List) value) { + terms.add(toLiteralTerms(o, codecRegistry)); + } + + return new ListTerm(terms); + } + + if (value instanceof Set) { + + List terms = new ArrayList<>(); + + for (Object o : (Set) value) { + terms.add(toLiteralTerms(o, codecRegistry)); + } + + return new SetTerm(terms); + } + + if (value instanceof Map) { + + Map terms = new LinkedHashMap<>(); + + ((Map) value).forEach((k, v) -> { + terms.put(toLiteralTerms(k, codecRegistry), toLiteralTerms(v, codecRegistry)); + }); + + return new MapTerm(terms); + } + + return QueryBuilder.literal(value, codecRegistry); + } + private SimpleStatementBuilder onBuild(SimpleStatementBuilder statementBuilder) { onBuild.forEach(it -> it.accept(statementBuilder)); @@ -264,4 +310,111 @@ public class StatementBuilder { */ BY_NAME; } + + static class ListTerm implements Term { + + private final Collection components; + + public ListTerm(@NonNull Collection components) { + this.components = components; + } + + @Override + public void appendTo(@NonNull StringBuilder builder) { + + if (components.isEmpty()) { + builder.append("[]"); + return; + } + + CqlHelper.append(components, builder, "[", ",", "]"); + } + + @Override + public boolean isIdempotent() { + for (Term component : components) { + if (!component.isIdempotent()) { + return false; + } + } + return true; + } + } + + static class SetTerm implements Term { + + private final Collection components; + + public SetTerm(@NonNull Collection components) { + this.components = components; + } + + @Override + public void appendTo(@NonNull StringBuilder builder) { + + if (components.isEmpty()) { + builder.append("{}"); + return; + } + + CqlHelper.append(components, builder, "{", ",", "}"); + } + + @Override + public boolean isIdempotent() { + for (Term component : components) { + if (!component.isIdempotent()) { + return false; + } + } + return true; + } + } + + static class MapTerm implements Term { + + private final Map components; + + public MapTerm(Map components) { + this.components = components; + } + + @Override + public void appendTo(@NonNull StringBuilder builder) { + + if (components.isEmpty()) { + builder.append("{}"); + return; + } + + boolean first = true; + + for (Map.Entry entry : components.entrySet()) { + + if (first) { + builder.append("{"); + first = false; + } else { + builder.append(","); + } + entry.getKey().appendTo(builder); + builder.append(":"); + entry.getValue().appendTo(builder); + } + if (!first) { + builder.append("}"); + } + } + + @Override + public boolean isIdempotent() { + for (Map.Entry entry : components.entrySet()) { + + if (!entry.getKey().isIdempotent() || !entry.getValue().isIdempotent()) { + return false; + } + } + return true; + } + } } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/mapping/BasicCassandraPersistentEntity.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/mapping/BasicCassandraPersistentEntity.java index e9888f92e..60afd41ed 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/mapping/BasicCassandraPersistentEntity.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/mapping/BasicCassandraPersistentEntity.java @@ -107,7 +107,7 @@ public class BasicCassandraPersistentEntity extends BasicPersistentEntity primaryKeyEntity = getRequiredPersistentEntity(property.getRawType()); for (CassandraPersistentProperty primaryKeyProperty : primaryKeyEntity) { + + DataType dataType = getDataTypeWithUserTypeFactory(primaryKeyProperty, DataTypeProvider.ShallowType); + if (primaryKeyProperty.isPartitionKeyColumn()) { - specification.partitionKeyColumn(primaryKeyProperty.getRequiredColumnName(), - getDataType(primaryKeyProperty)); + specification.partitionKeyColumn(primaryKeyProperty.getRequiredColumnName(), dataType); } else { // cluster column - specification.clusteredKeyColumn(primaryKeyProperty.getRequiredColumnName(), - getDataType(primaryKeyProperty), primaryKeyProperty.getPrimaryKeyOrdering()); + specification.clusteredKeyColumn(primaryKeyProperty.getRequiredColumnName(), dataType, + primaryKeyProperty.getPrimaryKeyOrdering()); } } } else { + DataType type = UserTypeUtil + .potentiallyFreeze(getDataTypeWithUserTypeFactory(property, DataTypeProvider.ShallowType)); + if (property.isIdProperty() || property.isPartitionKeyColumn()) { - specification.partitionKeyColumn(property.getRequiredColumnName(), - UserTypeUtil.potentiallyFreeze(getDataType(property))); + specification.partitionKeyColumn(property.getRequiredColumnName(), type); } else if (property.isClusterKeyColumn()) { - specification.clusteredKeyColumn(property.getRequiredColumnName(), - UserTypeUtil.potentiallyFreeze(getDataType(property)), property.getPrimaryKeyOrdering()); + specification.clusteredKeyColumn(property.getRequiredColumnName(), type, property.getPrimaryKeyOrdering()); } else { - specification.column(property.getRequiredColumnName(), UserTypeUtil.potentiallyFreeze(getDataType(property))); + specification.column(property.getRequiredColumnName(), type); } } } @@ -838,19 +841,31 @@ public class CassandraMappingContext } }, - FrozenLiteral { + ShallowType { @Override public DataType getDataType(CassandraPersistentEntity entity) { - return new ShallowUserDefinedType(com.datastax.oss.driver.api.core.CqlIdentifier.fromCql("system"), - entity.getTableName(), true); + return entity.isTupleType() ? entity.getTupleType() : new ShallowUserDefinedType(entity.getTableName(), false); } @Override DataType getUserType(com.datastax.oss.driver.api.core.CqlIdentifier userTypeName, UserTypeResolver userTypeResolver) { - return new ShallowUserDefinedType(com.datastax.oss.driver.api.core.CqlIdentifier.fromCql("system"), - userTypeName, true); + return new ShallowUserDefinedType(userTypeName, false); + } + }, + + FrozenLiteral { + + @Override + public DataType getDataType(CassandraPersistentEntity entity) { + return new ShallowUserDefinedType(entity.getTableName(), true); + } + + @Override + DataType getUserType(com.datastax.oss.driver.api.core.CqlIdentifier userTypeName, + UserTypeResolver userTypeResolver) { + return new ShallowUserDefinedType(userTypeName, true); } }; @@ -876,4 +891,113 @@ public class CassandraMappingContext UserTypeResolver userTypeResolver); } + + static class ShallowUserDefinedType implements com.datastax.oss.driver.api.core.type.UserDefinedType { + + private final CqlIdentifier name; + private final boolean frozen; + + public ShallowUserDefinedType(String name, boolean frozen) { + this(CqlIdentifier.fromInternal(name), frozen); + } + + public ShallowUserDefinedType(CqlIdentifier name, boolean frozen) { + this.name = name; + this.frozen = frozen; + } + + @Override + public CqlIdentifier getKeyspace() { + return null; + } + + @Override + public CqlIdentifier getName() { + return name; + } + + @Override + public boolean isFrozen() { + return frozen; + } + + @Override + public List getFieldNames() { + throw new UnsupportedOperationException( + "This implementation should only be used internally, this is likely a driver bug"); + } + + @Override + public int firstIndexOf(CqlIdentifier id) { + throw new UnsupportedOperationException( + "This implementation should only be used internally, this is likely a driver bug"); + } + + @Override + public int firstIndexOf(String name) { + throw new UnsupportedOperationException( + "This implementation should only be used internally, this is likely a driver bug"); + } + + @Override + public List getFieldTypes() { + throw new UnsupportedOperationException( + "This implementation should only be used internally, this is likely a driver bug"); + } + + @Override + public com.datastax.oss.driver.api.core.type.UserDefinedType copy(boolean newFrozen) { + return new ShallowUserDefinedType(this.name, newFrozen); + } + + @Override + public UdtValue newValue() { + throw new UnsupportedOperationException( + "This implementation should only be used internally, this is likely a driver bug"); + } + + @Override + public UdtValue newValue(@edu.umd.cs.findbugs.annotations.NonNull Object... fields) { + throw new UnsupportedOperationException( + "This implementation should only be used internally, this is likely a driver bug"); + } + + @Override + public AttachmentPoint getAttachmentPoint() { + throw new UnsupportedOperationException( + "This implementation should only be used internally, this is likely a driver bug"); + } + + @Override + public boolean isDetached() { + throw new UnsupportedOperationException( + "This implementation should only be used internally, this is likely a driver bug"); + } + + @Override + public void attach(@edu.umd.cs.findbugs.annotations.NonNull AttachmentPoint attachmentPoint) { + throw new UnsupportedOperationException( + "This implementation should only be used internally, this is likely a driver bug"); + } + + @Override + public boolean equals(Object o) { + if (this == o) + return true; + if (!(o instanceof com.datastax.oss.driver.api.core.type.UserDefinedType)) + return false; + com.datastax.oss.driver.api.core.type.UserDefinedType that = (com.datastax.oss.driver.api.core.type.UserDefinedType) o; + return isFrozen() == that.isFrozen() && Objects.equals(getName(), that.getName()); + } + + @Override + public int hashCode() { + return Objects.hash(name, frozen); + } + + @Override + public String toString() { + return "ShallowUserDefinedType{" + "name=" + name + ", frozen=" + frozen + '}'; + } + } } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/mapping/CassandraSimpleTypeHolder.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/mapping/CassandraSimpleTypeHolder.java index c6b1120a4..07ba570a6 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/mapping/CassandraSimpleTypeHolder.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/mapping/CassandraSimpleTypeHolder.java @@ -31,9 +31,9 @@ import org.springframework.lang.Nullable; import com.datastax.oss.driver.api.core.cql.Row; import com.datastax.oss.driver.api.core.data.TupleValue; +import com.datastax.oss.driver.api.core.data.UdtValue; import com.datastax.oss.driver.api.core.type.DataType; import com.datastax.oss.driver.api.core.type.DataTypes; -import com.datastax.oss.driver.api.core.type.UserDefinedType; import com.datastax.oss.driver.api.core.type.codec.TypeCodec; import com.datastax.oss.driver.api.core.type.codec.registry.CodecRegistry; import com.datastax.oss.driver.api.core.type.reflect.GenericType; @@ -82,7 +82,7 @@ public class CassandraSimpleTypeHolder extends SimpleTypeHolder { simpleTypes.add(Number.class); simpleTypes.add(Row.class); simpleTypes.add(TupleValue.class); - simpleTypes.add(UserDefinedType.class); + simpleTypes.add(UdtValue.class); classToDataType = Collections.unmodifiableMap(classToDataType(codecRegistry, primitiveWrappers)); nameToDataType = Collections.unmodifiableMap(nameToDataType()); diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/mapping/UserTypeUtil.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/mapping/UserTypeUtil.java index 1d2b45d8f..d1077257f 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/mapping/UserTypeUtil.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/mapping/UserTypeUtil.java @@ -23,7 +23,6 @@ import com.datastax.oss.driver.api.core.type.ListType; import com.datastax.oss.driver.api.core.type.MapType; import com.datastax.oss.driver.api.core.type.SetType; import com.datastax.oss.driver.api.core.type.UserDefinedType; -import com.datastax.oss.driver.internal.core.metadata.schema.ShallowUserDefinedType; /** * {@link com.datastax.driver.core.UserType} utility methods. Mainly for internal use within the framework. @@ -78,8 +77,7 @@ class UserTypeUtil { } if (isNonFrozenUdt(dataType)) { - UserDefinedType userDefinedType = (UserDefinedType) dataType; - return new ShallowUserDefinedType(userDefinedType.getKeyspace(), userDefinedType.getName(), true); + return ((UserDefinedType) dataType).copy(true); } return dataType; @@ -92,5 +90,4 @@ class UserTypeUtil { private static boolean isNonFrozenUdt(DataType dataType) { return dataType instanceof UserDefinedType && !((UserDefinedType) dataType).isFrozen(); } - } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/query/CassandraPageRequest.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/query/CassandraPageRequest.java index deadf7ea9..f864ca669 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/query/CassandraPageRequest.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/query/CassandraPageRequest.java @@ -177,7 +177,12 @@ public class CassandraPageRequest extends PageRequest { */ @Nullable public ByteBuffer getPagingState() { - return this.pagingState; + + if (this.pagingState == null) { + return null; + } + + return this.pagingState.asReadOnlyBuffer(); } /** @@ -198,7 +203,7 @@ public class CassandraPageRequest extends PageRequest { Assert.state(hasNext(), "Cannot create a next page request without a PagingState"); - return new CassandraPageRequest(getPageNumber() + 1, getPageSize(), getSort(), this.pagingState, false); + return new CassandraPageRequest(getPageNumber() + 1, getPageSize(), getSort(), getPagingState(), false); } /** @@ -212,7 +217,7 @@ public class CassandraPageRequest extends PageRequest { Assert.notNull(sort, "Sort must not be null"); - return new CassandraPageRequest(this.getPageNumber(), this.getPageSize(), sort, this.pagingState, this.nextAllowed); + return new CassandraPageRequest(this.getPageNumber(), this.getPageSize(), sort, getPagingState(), this.nextAllowed); } /* (non-Javadoc) diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/query/CriteriaDefinition.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/query/CriteriaDefinition.java index 4f70cb095..15993628d 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/query/CriteriaDefinition.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/query/CriteriaDefinition.java @@ -18,6 +18,7 @@ package org.springframework.data.cassandra.core.query; import lombok.EqualsAndHashCode; import java.util.Optional; +import java.util.function.Function; import org.springframework.lang.Nullable; import org.springframework.util.Assert; @@ -84,6 +85,18 @@ public interface CriteriaDefinition { public Object getValue() { return this.value; } + + /** + * This method allows the application of a function to this {@link Predicate} value. The function should expect a + * single {@link Object} argument and produce an {@code R} result. Any exception thrown by f() will be propagated to + * the caller. + * + * @param + * @return the result of the {@link Function mappingFunction}. + */ + public R as(Function mappingFunction) { + return mappingFunction.apply(this.value); + } } /** diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/config/KeyspaceActionSpecificationFactoryBeanUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/config/KeyspaceActionSpecificationFactoryBeanUnitTests.java index 44bd59baa..b2c80d52e 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/config/KeyspaceActionSpecificationFactoryBeanUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/config/KeyspaceActionSpecificationFactoryBeanUnitTests.java @@ -21,13 +21,15 @@ import java.util.Arrays; import java.util.List; import org.junit.Test; -import org.springframework.data.cassandra.core.cql.KeyspaceIdentifier; + import org.springframework.data.cassandra.core.cql.keyspace.AlterKeyspaceSpecification; import org.springframework.data.cassandra.core.cql.keyspace.CreateKeyspaceSpecification; import org.springframework.data.cassandra.core.cql.keyspace.DropKeyspaceSpecification; import org.springframework.data.cassandra.core.cql.keyspace.KeyspaceActionSpecification; import org.springframework.data.cassandra.core.cql.keyspace.KeyspaceOption.ReplicationStrategy; +import com.datastax.oss.driver.api.core.CqlIdentifier; + /** * Unit tests for {@link KeyspaceActionSpecificationFactoryBean}. * @@ -53,7 +55,7 @@ public class KeyspaceActionSpecificationFactoryBeanUnitTests { CreateKeyspaceSpecification create = (CreateKeyspaceSpecification) actions.get(0); - assertThat(create.getName()).isEqualTo(KeyspaceIdentifier.of("my_keyspace")); + assertThat(create.getName()).isEqualTo(CqlIdentifier.fromCql("my_keyspace")); assertThat(create.getOptions()).containsKeys("durable_writes", "replication"); } @@ -72,7 +74,7 @@ public class KeyspaceActionSpecificationFactoryBeanUnitTests { DropKeyspaceSpecification drop = (DropKeyspaceSpecification) actions.get(1); - assertThat(drop.getName()).isEqualTo(KeyspaceIdentifier.of("my_keyspace")); + assertThat(drop.getName()).isEqualTo(CqlIdentifier.fromCql("my_keyspace")); } @Test // DATACASS-502 @@ -93,7 +95,7 @@ public class KeyspaceActionSpecificationFactoryBeanUnitTests { AlterKeyspaceSpecification alter = (AlterKeyspaceSpecification) actions.get(0); - assertThat(alter.getName()).isEqualTo(KeyspaceIdentifier.of("my_keyspace")); + assertThat(alter.getName()).isEqualTo(CqlIdentifier.fromCql("my_keyspace")); assertThat(alter.getOptions()).containsKeys("durable_writes", "replication"); } @@ -113,7 +115,7 @@ public class KeyspaceActionSpecificationFactoryBeanUnitTests { AlterKeyspaceSpecification alter = (AlterKeyspaceSpecification) actions.get(0); - assertThat(alter.getName()).isEqualTo(KeyspaceIdentifier.of("my_keyspace")); + assertThat(alter.getName()).isEqualTo(CqlIdentifier.fromCql("my_keyspace")); assertThat(alter.getOptions()).containsKeys("durable_writes", "replication"); } @@ -131,7 +133,7 @@ public class KeyspaceActionSpecificationFactoryBeanUnitTests { AlterKeyspaceSpecification alter = (AlterKeyspaceSpecification) actions.get(0); - assertThat(alter.getName()).isEqualTo(KeyspaceIdentifier.of("my_keyspace")); + assertThat(alter.getName()).isEqualTo(CqlIdentifier.fromCql("my_keyspace")); assertThat(alter.getOptions()).doesNotContainKeys("replication"); } } diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/config/SchemaActionIntegrationTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/config/SchemaActionIntegrationTests.java index b10a4c665..316c51781 100755 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/config/SchemaActionIntegrationTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/config/SchemaActionIntegrationTests.java @@ -19,44 +19,42 @@ package org.springframework.data.cassandra.config; import static org.assertj.core.api.Assertions.*; import java.util.Collections; -import java.util.List; import java.util.Set; -import org.junit.Before; import org.junit.Test; import org.springframework.beans.factory.BeanCreationException; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.AnnotationConfigApplicationContext; import org.springframework.context.annotation.Configuration; +import org.springframework.core.io.ByteArrayResource; import org.springframework.data.cassandra.core.cql.SessionCallback; +import org.springframework.data.cassandra.core.cql.session.init.KeyspacePopulator; +import org.springframework.data.cassandra.core.cql.session.init.ResourceKeyspacePopulator; import org.springframework.data.cassandra.domain.Person; -import org.springframework.data.cassandra.test.util.AbstractKeyspaceCreatingIntegrationTest; +import org.springframework.data.cassandra.repository.support.IntegrationTestConfig; +import org.springframework.data.cassandra.test.util.AbstractEmbeddedCassandraIntegrationTest; +import org.springframework.lang.Nullable; import com.datastax.oss.driver.api.core.CqlSession; import com.datastax.oss.driver.api.core.metadata.schema.KeyspaceMetadata; import com.datastax.oss.driver.api.core.metadata.schema.TableMetadata; /** - * Integration test testing various {@link SchemaAction SchemaActions} on startup of a Spring configured, - * Apache Cassandra application client. + * Integration test testing various {@link SchemaAction SchemaActions} on startup of a Spring configured, Apache + * Cassandra application client. * * @author John Blum * @author Mark Paluch * @see org.springframework.data.cassandra.test.util.AbstractKeyspaceCreatingIntegrationTest */ -public class SchemaActionIntegrationTests extends AbstractKeyspaceCreatingIntegrationTest { +public class SchemaActionIntegrationTests extends AbstractEmbeddedCassandraIntegrationTest { - protected static final String CREATE_PERSON_TABLE_CQL = - "CREATE TABLE IF NOT EXISTS person (id int, firstName text, lastName text, PRIMARY KEY(id));"; - - protected static final String DROP_ADDRESS_TYPE_CQL = "DROP TYPE IF EXISTS address"; - protected static final String DROP_PERSON_TABLE_CQL = "DROP TABLE IF EXISTS person"; + protected static final String CREATE_PERSON_TABLE_CQL = "CREATE TABLE IF NOT EXISTS person (id int, firstName text, lastName text, PRIMARY KEY(id));"; protected ConfigurableApplicationContext newApplicationContext(Class... annotatedClasses) { - AnnotationConfigApplicationContext applicationContext = - new AnnotationConfigApplicationContext(annotatedClasses); + AnnotationConfigApplicationContext applicationContext = new AnnotationConfigApplicationContext(annotatedClasses); applicationContext.registerShutdownHook(); @@ -74,7 +72,7 @@ public class SchemaActionIntegrationTests extends AbstractKeyspaceCreatingIntegr @SuppressWarnings("all") protected void assertTableWithColumnsExists(CqlSession session, String tableName, String... columns) { - KeyspaceMetadata keyspaceMetadata = session.getMetadata().getKeyspace(session.getKeyspace().get()).orElse(null); + KeyspaceMetadata keyspaceMetadata = session.refreshSchema().getKeyspace(session.getKeyspace().get()).orElse(null); assertThat(keyspaceMetadata).isNotNull(); @@ -90,23 +88,13 @@ public class SchemaActionIntegrationTests extends AbstractKeyspaceCreatingIntegr assertThat(tableMetadata.getColumns()).hasSize(columns.length); } - @Before - public void setup() { - - CqlSession session = getSession(); - - session.execute(DROP_PERSON_TABLE_CQL); - session.execute(DROP_ADDRESS_TYPE_CQL); - } - @Test public void createWithNoExistingTableCreatesTableFromEntity() { doInSessionWithConfiguration(CreateWithNoExistingTableConfiguration.class, session -> { - assertTableWithColumnsExists(session, "person", "firstName", "lastName", "nickname", - "birthDate", "numberOfChildren", "cool", "createdDate", "zoneId", "mainAddress", - "alternativeAddresses"); + assertTableWithColumnsExists(session, "person", "firstName", "lastName", "nickname", "birthDate", + "numberOfChildren", "cool", "createdDate", "zoneId", "mainAddress", "alternativeAddresses"); return null; }); @@ -118,6 +106,7 @@ public class SchemaActionIntegrationTests extends AbstractKeyspaceCreatingIntegr try { doInSessionWithConfiguration(CreateWithExistingTableConfiguration.class, session -> { + fail(String.format("%s should have failed", CreateWithExistingTableConfiguration.class.getSimpleName())); return null; }); @@ -125,7 +114,7 @@ public class SchemaActionIntegrationTests extends AbstractKeyspaceCreatingIntegr fail("Expected BeanCreationException"); } catch (BeanCreationException cause) { - assertThat(cause).hasMessageContaining(String.format("Table %s.person already exists", getKeyspace())); + assertThat(cause).hasMessageContaining("person already exists"); } } @@ -134,9 +123,8 @@ public class SchemaActionIntegrationTests extends AbstractKeyspaceCreatingIntegr doInSessionWithConfiguration(CreateIfNotExistsWithNoExistingTableConfiguration.class, session -> { - assertTableWithColumnsExists(session, "person", "firstName", "lastName", "nickname", - "birthDate", "numberOfChildren", "cool", "createdDate", "zoneId", "mainAddress", - "alternativeAddresses"); + assertTableWithColumnsExists(session, "person", "firstName", "lastName", "nickname", "birthDate", + "numberOfChildren", "cool", "createdDate", "zoneId", "mainAddress", "alternativeAddresses"); return null; }); @@ -158,25 +146,15 @@ public class SchemaActionIntegrationTests extends AbstractKeyspaceCreatingIntegr doInSessionWithConfiguration(RecreateSchemaActionWithExistingTableConfiguration.class, session -> { - assertTableWithColumnsExists(session, "person", "firstName", "lastName", "nickname", - "birthDate", "numberOfChildren", "cool", "createdDate", "zoneId", "mainAddress", - "alternativeAddresses"); + assertTableWithColumnsExists(session, "person", "firstName", "lastName", "nickname", "birthDate", + "numberOfChildren", "cool", "createdDate", "zoneId", "mainAddress", "alternativeAddresses"); return null; }); } @Configuration - static class CreateWithNoExistingTableConfiguration extends CassandraConfiguration { - - @Override - public SchemaAction getSchemaAction() { - return SchemaAction.CREATE; - } - } - - @Configuration - static class CreateWithExistingTableConfiguration extends CassandraConfiguration { + static class CreateWithNoExistingTableConfiguration extends IntegrationTestConfig { @Override public SchemaAction getSchemaAction() { @@ -184,22 +162,33 @@ public class SchemaActionIntegrationTests extends AbstractKeyspaceCreatingIntegr } @Override - protected List getStartupScripts() { - return Collections.singletonList(CREATE_PERSON_TABLE_CQL); + protected Set> getInitialEntitySet() throws ClassNotFoundException { + return Collections.singleton(Person.class); } } @Configuration - static class CreateIfNotExistsWithNoExistingTableConfiguration extends CassandraConfiguration { + static class CreateWithExistingTableConfiguration extends IntegrationTestConfig { @Override public SchemaAction getSchemaAction() { - return SchemaAction.CREATE_IF_NOT_EXISTS; + return SchemaAction.CREATE; + } + + @Nullable + @Override + protected KeyspacePopulator keyspacePopulator() { + return new ResourceKeyspacePopulator(new ByteArrayResource(CREATE_PERSON_TABLE_CQL.getBytes())); + } + + @Override + protected Set> getInitialEntitySet() throws ClassNotFoundException { + return Collections.singleton(Person.class); } } @Configuration - static class CreateIfNotExistsWithExistingTableConfiguration extends CassandraConfiguration { + static class CreateIfNotExistsWithNoExistingTableConfiguration extends IntegrationTestConfig { @Override public SchemaAction getSchemaAction() { @@ -207,36 +196,48 @@ public class SchemaActionIntegrationTests extends AbstractKeyspaceCreatingIntegr } @Override - protected List getStartupScripts() { - return Collections.singletonList(CREATE_PERSON_TABLE_CQL); + protected Set> getInitialEntitySet() throws ClassNotFoundException { + return Collections.singleton(Person.class); } } @Configuration - static class RecreateSchemaActionWithExistingTableConfiguration extends CassandraConfiguration { + static class CreateIfNotExistsWithExistingTableConfiguration extends IntegrationTestConfig { + + @Override + public SchemaAction getSchemaAction() { + return SchemaAction.CREATE_IF_NOT_EXISTS; + } + + @Nullable + @Override + protected KeyspacePopulator keyspacePopulator() { + return new ResourceKeyspacePopulator(new ByteArrayResource(CREATE_PERSON_TABLE_CQL.getBytes())); + } + + @Override + protected Set> getInitialEntitySet() throws ClassNotFoundException { + return Collections.singleton(Person.class); + } + } + + @Configuration + static class RecreateSchemaActionWithExistingTableConfiguration extends IntegrationTestConfig { @Override public SchemaAction getSchemaAction() { return SchemaAction.RECREATE; } + @Nullable @Override - protected List getStartupScripts() { - return Collections.singletonList(CREATE_PERSON_TABLE_CQL); + protected KeyspacePopulator keyspacePopulator() { + return new ResourceKeyspacePopulator(new ByteArrayResource(CREATE_PERSON_TABLE_CQL.getBytes())); } - } - - @Configuration - static abstract class CassandraConfiguration extends AbstractCassandraConfiguration { @Override - protected Set> getInitialEntitySet() { + protected Set> getInitialEntitySet() throws ClassNotFoundException { return Collections.singleton(Person.class); } - - @Override - protected String getKeyspaceName() { - return keyspaceRule.getKeyspaceName(); - } } } diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplateIntegrationTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplateIntegrationTests.java index b0beffa0c..7997d2543 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplateIntegrationTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/AsyncCassandraTemplateIntegrationTests.java @@ -204,7 +204,7 @@ public class AsyncCassandraTemplateIntegrationTests extends AbstractKeyspaceCrea } @Test // DATACASS-292 - public void updateShouldUpdateEntityWithLwt() { + public void updateShouldUpdateEntityWithLwt() throws InterruptedException { UpdateOptions lwtOptions = UpdateOptions.builder().withIfExists().build(); @@ -217,6 +217,10 @@ public class AsyncCassandraTemplateIntegrationTests extends AbstractKeyspaceCrea assertThat(getUninterruptibly(updated).wasApplied()).isTrue(); assertThat(getUninterruptibly(updated).getEntity()).isSameAs(user); + + // Cassandra requires a while to apply that change... + Thread.sleep(200); + assertThat(getUser(user.getId()).getFirstname()).isEqualTo("Walter Hartwell"); } 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 9b9525afb..5a1c5f2f7 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 @@ -46,12 +46,14 @@ import org.springframework.data.cassandra.domain.VersionedUser; import org.springframework.data.mapping.callback.EntityCallbacks; 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.cql.AsyncResultSet; 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; +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; @@ -69,7 +71,7 @@ public class AsyncCassandraTemplateUnitTests { @Mock ColumnDefinition columnDefinition; @Mock ColumnDefinitions columnDefinitions; - @Captor ArgumentCaptor statementCaptor; + @Captor ArgumentCaptor statementCaptor; AsyncCassandraTemplate template; @@ -108,7 +110,7 @@ public class AsyncCassandraTemplateUnitTests { public void selectUsingCqlShouldReturnMappedResults() { when(resultSet.currentPage()).thenReturn(Collections.singleton(row)); - when(columnDefinitions.contains(anyString())).thenReturn(true); + when(columnDefinitions.contains(any(CqlIdentifier.class))).thenReturn(true); when(columnDefinitions.get(anyInt())).thenReturn(columnDefinition); when(columnDefinitions.firstIndexOf("id")).thenReturn(0); @@ -125,15 +127,14 @@ public class AsyncCassandraTemplateUnitTests { assertThat(getUninterruptibly(list)).hasSize(1).contains(new User("myid", "Walter", "White")); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users"); } @Test // DATACASS-292 public void selectUsingCqlShouldInvokeCallbackWithMappedResults() { when(resultSet.currentPage()).thenReturn(Collections.singletonList(row)); - when(columnDefinitions.contains(anyString())).thenReturn(true); - + when(columnDefinitions.contains(any(CqlIdentifier.class))).thenReturn(true); when(columnDefinitions.get(anyInt())).thenReturn(columnDefinition); when(columnDefinitions.firstIndexOf("id")).thenReturn(0); when(columnDefinitions.firstIndexOf("firstname")).thenReturn(1); @@ -152,7 +153,7 @@ public class AsyncCassandraTemplateUnitTests { assertThat(getUninterruptibly(result)).isNull(); assertThat(list).hasSize(1).contains(new User("myid", "Walter", "White")); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users"); } @Test // DATACASS-292 @@ -176,7 +177,7 @@ public class AsyncCassandraTemplateUnitTests { public void selectOneShouldReturnMappedResults() { when(resultSet.currentPage()).thenReturn(Collections.singleton(row)); - when(columnDefinitions.contains(anyString())).thenReturn(true); + when(columnDefinitions.contains(any(CqlIdentifier.class))).thenReturn(true); when(columnDefinitions.get(anyInt())).thenReturn(columnDefinition); when(columnDefinitions.firstIndexOf("id")).thenReturn(0); @@ -189,18 +190,18 @@ public class AsyncCassandraTemplateUnitTests { when(row.getObject(1)).thenReturn("Walter"); when(row.getObject(2)).thenReturn("White"); - ListenableFuture future = template.selectOne("SELECT * FROM users WHERE id='myid';", User.class); + ListenableFuture future = template.selectOne("SELECT * FROM users WHERE id='myid'", User.class); assertThat(getUninterruptibly(future)).isEqualTo(new User("myid", "Walter", "White")); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users WHERE id='myid';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users WHERE id='myid'"); } @Test // DATACASS-292 public void selectOneByIdShouldReturnMappedResults() { when(resultSet.currentPage()).thenReturn(Collections.singleton(row)); - when(columnDefinitions.contains(anyString())).thenReturn(true); + when(columnDefinitions.contains(any(CqlIdentifier.class))).thenReturn(true); when(columnDefinitions.get(anyInt())).thenReturn(columnDefinition); when(columnDefinitions.firstIndexOf("id")).thenReturn(0); when(columnDefinitions.firstIndexOf("firstname")).thenReturn(1); @@ -216,7 +217,7 @@ public class AsyncCassandraTemplateUnitTests { assertThat(getUninterruptibly(future)).isEqualTo(new User("myid", "Walter", "White")); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users WHERE id='myid';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users WHERE id='myid' LIMIT 1"); } @Test // DATACASS-696 @@ -224,7 +225,7 @@ public class AsyncCassandraTemplateUnitTests { when(resultSet.currentPage()).thenReturn(Collections.singleton(row)); - ListenableFuture future = template.selectOne("SELECT id FROM users WHERE id='myid';", String.class); + ListenableFuture future = template.selectOne("SELECT id FROM users WHERE id='myid'", String.class); assertThat(getUninterruptibly(future)).isNull(); } @@ -232,37 +233,35 @@ public class AsyncCassandraTemplateUnitTests { @Test // DATACASS-292 public void existsShouldReturnExistingElement() { - when(resultSet.currentPage()).thenReturn(Collections.singleton(row)); + when(resultSet.one()).thenReturn(row); ListenableFuture future = template.exists("myid", User.class); assertThat(getUninterruptibly(future)).isTrue(); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users WHERE id='myid';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users WHERE id='myid' LIMIT 1"); } @Test // DATACASS-292 public void existsShouldReturnNonExistingElement() { - when(resultSet.currentPage()).thenReturn(Collections.emptyList()); - ListenableFuture future = template.exists("myid", User.class); assertThat(getUninterruptibly(future)).isFalse(); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users WHERE id='myid';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users WHERE id='myid' LIMIT 1"); } @Test // DATACASS-512 public void existsByQueryShouldReturnExistingElement() { - when(resultSet.currentPage()).thenReturn(Collections.singleton(row)); + when(resultSet.one()).thenReturn(row); ListenableFuture future = template.exists(Query.empty(), User.class); assertThat(getUninterruptibly(future)).isTrue(); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users LIMIT 1;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users LIMIT 1"); } @Test // DATACASS-292 @@ -276,7 +275,7 @@ public class AsyncCassandraTemplateUnitTests { assertThat(getUninterruptibly(future)).isEqualTo(42L); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT count(*) FROM users;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT count(1) FROM users"); } @Test // DATACASS-292 @@ -290,7 +289,7 @@ public class AsyncCassandraTemplateUnitTests { assertThat(getUninterruptibly(future)).isEqualTo(42L); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT COUNT(1) FROM users;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT count(1) FROM users"); } @Test // DATACASS-292, DATACASS-618 @@ -304,8 +303,8 @@ public class AsyncCassandraTemplateUnitTests { assertThat(getUninterruptibly(future)).isEqualTo(user); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()) - .hasToString("INSERT INTO users (firstname,id,lastname) VALUES ('Walter','heisenberg','White');"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("INSERT INTO users (firstname,id,lastname) VALUES ('Walter','heisenberg','White')"); assertThat(beforeConvert).isSameAs(user); assertThat(beforeSave).isSameAs(user); } @@ -321,8 +320,8 @@ public class AsyncCassandraTemplateUnitTests { assertThat(getUninterruptibly(future)).isEqualTo(user); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString( - "INSERT INTO vusers (firstname,id,lastname,version) VALUES ('Walter','heisenberg','White',0) IF NOT EXISTS;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo( + "INSERT INTO vusers (firstname,id,lastname,version) VALUES ('Walter','heisenberg','White',0) IF NOT EXISTS"); assertThat(beforeConvert).isSameAs(user); assertThat(beforeSave).isSameAs(user); } @@ -357,8 +356,8 @@ public class AsyncCassandraTemplateUnitTests { assertThat(getUninterruptibly(future)).isEqualTo(user); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()) - .hasToString("UPDATE users SET firstname='Walter',lastname='White' WHERE id='heisenberg';"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("UPDATE users SET firstname='Walter', lastname='White' WHERE id='heisenberg'"); assertThat(beforeConvert).isSameAs(user); assertThat(beforeSave).isSameAs(user); } @@ -375,8 +374,8 @@ public class AsyncCassandraTemplateUnitTests { assertThat(getUninterruptibly(future)).isEqualTo(user); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString( - "UPDATE vusers SET firstname='Walter',lastname='White',version=1 WHERE id='heisenberg' IF version=0;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo( + "UPDATE vusers SET firstname='Walter', lastname='White', version=1 WHERE id='heisenberg' IF version=0"); assertThat(beforeConvert).isSameAs(user); assertThat(beforeSave).isSameAs(user); } @@ -390,8 +389,8 @@ public class AsyncCassandraTemplateUnitTests { template.update(user, updateOptions); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()) - .hasToString("UPDATE users SET firstname='Walter',lastname='White' WHERE id='heisenberg' IF EXISTS;"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("UPDATE users SET firstname='Walter', lastname='White' WHERE id='heisenberg' IF EXISTS"); } @Test // DATACASS-575 @@ -403,8 +402,8 @@ public class AsyncCassandraTemplateUnitTests { template.update(user, options); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString( - "UPDATE users SET firstname='Walter',lastname='White' WHERE id='heisenberg' IF firstname='Walter';"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("UPDATE users SET firstname='Walter', lastname='White' WHERE id='heisenberg' IF firstname='Walter'"); } @Test // DATACASS-575 @@ -416,7 +415,8 @@ public class AsyncCassandraTemplateUnitTests { template.update(query, update, User.class); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("UPDATE users SET firstname='Walter' WHERE id='heisenberg';"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("UPDATE users SET firstname='Walter' WHERE id='heisenberg'"); } @Test // DATACASS-575 @@ -432,8 +432,8 @@ public class AsyncCassandraTemplateUnitTests { template.update(query, update, User.class); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString( - "UPDATE users SET firstname='Walter' WHERE id='heisenberg' IF firstname='Walter' AND lastname='White';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo( + "UPDATE users SET firstname='Walter' WHERE id='heisenberg' IF firstname='Walter' AND lastname='White'"); } @Test // DATACASS-292 @@ -466,7 +466,7 @@ public class AsyncCassandraTemplateUnitTests { assertThat(getUninterruptibly(future)).isTrue(); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("DELETE FROM users WHERE id='heisenberg';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("DELETE FROM users WHERE id='heisenberg'"); } @Test // DATACASS-292 @@ -480,7 +480,7 @@ public class AsyncCassandraTemplateUnitTests { assertThat(getUninterruptibly(future)).isEqualTo(user); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("DELETE FROM users WHERE id='heisenberg';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("DELETE FROM users WHERE id='heisenberg'"); } @Test // DATACASS-575 @@ -492,8 +492,8 @@ public class AsyncCassandraTemplateUnitTests { template.delete(user, options); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()) - .hasToString("DELETE FROM users WHERE id='heisenberg' IF firstname='Walter';"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("DELETE FROM users WHERE id='heisenberg' IF firstname='Walter'"); } @Test // DATACASS-575 @@ -505,8 +505,8 @@ public class AsyncCassandraTemplateUnitTests { template.delete(query, User.class); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()) - .hasToString("DELETE FROM users WHERE id='heisenberg' IF firstname='Walter';"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("DELETE FROM users WHERE id='heisenberg' IF firstname='Walter'"); } @Test // DATACASS-292 @@ -534,7 +534,7 @@ public class AsyncCassandraTemplateUnitTests { template.truncate(User.class); verify(session).executeAsync(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("TRUNCATE users;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("TRUNCATE users"); } private static T getUninterruptibly(Future future) { diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraPersistentEntitySchemaDropperUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraPersistentEntitySchemaDropperUnitTests.java index c1ea83077..d9feb80ca 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraPersistentEntitySchemaDropperUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraPersistentEntitySchemaDropperUnitTests.java @@ -110,7 +110,8 @@ public class CassandraPersistentEntitySchemaDropperUnitTests extends CassandraPe context.setInitialEntitySet(Collections.singleton(Person.class)); context.afterPropertiesSet(); - when(metadata.getTables()).thenReturn(createTables(person, contact)); + Map tables = createTables(person, contact); + when(metadata.getTables()).thenReturn(tables); CassandraPersistentEntitySchemaDropper schemaDropper = new CassandraPersistentEntitySchemaDropper(context, operations); @@ -129,7 +130,8 @@ public class CassandraPersistentEntitySchemaDropperUnitTests extends CassandraPe context.setInitialEntitySet(Collections.singleton(Person.class)); context.afterPropertiesSet(); - when(metadata.getTables()).thenReturn(createTables(person, contact)); + Map tables = createTables(person, contact); + when(metadata.getTables()).thenReturn(tables); CassandraPersistentEntitySchemaDropper schemaDropper = new CassandraPersistentEntitySchemaDropper(context, operations); diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateIntegrationTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateIntegrationTests.java index 124288b81..0d319d994 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateIntegrationTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateIntegrationTests.java @@ -316,7 +316,7 @@ public class CassandraTemplateIntegrationTests extends AbstractKeyspaceCreatingI } @Test // DATACASS-292 - public void updateShouldUpdateEntityWithLwt() { + public void updateShouldUpdateEntityWithLwt() throws InterruptedException { UpdateOptions lwtOptions = UpdateOptions.builder().withIfExists().build(); @@ -329,6 +329,9 @@ public class CassandraTemplateIntegrationTests extends AbstractKeyspaceCreatingI WriteResult lwt = template.update(user, lwtOptions); assertThat(lwt.wasApplied()).isTrue(); + + // Await until Cassandra has persisted the change + Thread.sleep(300); assertThat(template.selectOneById(user.getId(), User.class).getFirstname()).isEqualTo("Walter Hartwell"); } 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 39162ec29..b952fbefc 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 @@ -24,7 +24,6 @@ import java.util.Collections; import java.util.List; import org.junit.Before; -import org.junit.Ignore; import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.ArgumentCaptor; @@ -42,12 +41,14 @@ import org.springframework.data.cassandra.domain.User; import org.springframework.data.cassandra.domain.VersionedUser; 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.cql.ColumnDefinition; import com.datastax.oss.driver.api.core.cql.ColumnDefinitions; import com.datastax.oss.driver.api.core.cql.ResultSet; 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.type.DataTypes; @@ -65,7 +66,7 @@ public class CassandraTemplateUnitTests { @Mock ColumnDefinition columnDefinition; @Mock ColumnDefinitions columnDefinitions; - @Captor ArgumentCaptor> statementCaptor; + @Captor ArgumentCaptor statementCaptor; CassandraTemplate template; @@ -104,7 +105,7 @@ public class CassandraTemplateUnitTests { public void selectUsingCqlShouldReturnMappedResults() { when(resultSet.iterator()).thenReturn(Collections.singleton(row).iterator()); - when(columnDefinitions.contains(anyString())).thenReturn(true); + when(columnDefinitions.contains(any(CqlIdentifier.class))).thenReturn(true); when(columnDefinitions.get(anyInt())).thenReturn(columnDefinition); when(columnDefinitions.firstIndexOf("id")).thenReturn(0); @@ -121,7 +122,7 @@ public class CassandraTemplateUnitTests { assertThat(list).hasSize(1).contains(new User("myid", "Walter", "White")); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users"); } @Test // DATACASS-292 @@ -142,7 +143,7 @@ public class CassandraTemplateUnitTests { public void selectOneShouldReturnMappedResults() { when(resultSet.iterator()).thenReturn(Collections.singleton(row).iterator()); - when(columnDefinitions.contains(anyString())).thenReturn(true); + when(columnDefinitions.contains(any(CqlIdentifier.class))).thenReturn(true); when(columnDefinitions.get(anyInt())).thenReturn(columnDefinition); when(columnDefinitions.firstIndexOf("id")).thenReturn(0); when(columnDefinitions.firstIndexOf("firstname")).thenReturn(1); @@ -154,11 +155,11 @@ public class CassandraTemplateUnitTests { when(row.getObject(1)).thenReturn("Walter"); when(row.getObject(2)).thenReturn("White"); - User user = template.selectOne("SELECT * FROM users WHERE id='myid';", User.class); + User user = template.selectOne("SELECT * FROM users WHERE id='myid'", User.class); assertThat(user).isEqualTo(new User("myid", "Walter", "White")); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users WHERE id='myid';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users WHERE id='myid'"); } @Test // DATACASS-696 @@ -166,7 +167,7 @@ public class CassandraTemplateUnitTests { when(resultSet.iterator()).thenReturn(Collections.singleton(row).iterator()); - String nullValue = template.selectOne("SELECT id FROM users WHERE id='myid';", String.class); + String nullValue = template.selectOne("SELECT id FROM users WHERE id='myid'", String.class); assertThat(nullValue).isNull(); } @@ -175,7 +176,7 @@ public class CassandraTemplateUnitTests { public void selectOneByIdShouldReturnMappedResults() { when(resultSet.iterator()).thenReturn(Collections.singleton(row).iterator()); - when(columnDefinitions.contains(anyString())).thenReturn(true); + when(columnDefinitions.contains(any(CqlIdentifier.class))).thenReturn(true); when(columnDefinitions.get(anyInt())).thenReturn(columnDefinition); when(columnDefinitions.firstIndexOf("id")).thenReturn(0); when(columnDefinitions.firstIndexOf("firstname")).thenReturn(1); @@ -191,14 +192,14 @@ public class CassandraTemplateUnitTests { assertThat(user).isEqualTo(new User("myid", "Walter", "White")); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users WHERE id='myid';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users WHERE id='myid' LIMIT 1"); } @Test // DATACASS-313 public void selectProjectedOneShouldReturnMappedResults() { when(resultSet.iterator()).thenReturn(Collections.singleton(row).iterator()); - when(columnDefinitions.contains(anyString())).thenReturn(true); + when(columnDefinitions.contains(any(CqlIdentifier.class))).thenReturn(true); when(columnDefinitions.get(anyInt())).thenReturn(columnDefinition); when(columnDefinitions.firstIndexOf("id")).thenReturn(0); @@ -210,43 +211,41 @@ public class CassandraTemplateUnitTests { assertThat(user.getFirstname()).isEqualTo("Walter"); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT firstname FROM users LIMIT 2;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT firstname FROM users LIMIT 2"); } @Test // DATACASS-292 public void existsShouldReturnExistingElement() { - when(resultSet.iterator()).thenReturn(Collections.singleton(row).iterator()); + when(resultSet.one()).thenReturn(row); boolean exists = template.exists("myid", User.class); assertThat(exists).isTrue(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users WHERE id='myid';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users WHERE id='myid' LIMIT 1"); } @Test // DATACASS-292 public void existsShouldReturnNonExistingElement() { - when(resultSet.iterator()).thenReturn(Collections.emptyIterator()); - boolean exists = template.exists("myid", User.class); assertThat(exists).isFalse(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users WHERE id='myid';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users WHERE id='myid' LIMIT 1"); } @Test // DATACASS-512 public void existsByQueryShouldReturnExistingElement() { - when(resultSet.iterator()).thenReturn(Collections.singleton(row).iterator()); + when(resultSet.one()).thenReturn(row); boolean exists = template.exists(Query.empty(), User.class); assertThat(exists).isTrue(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users LIMIT 1;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users LIMIT 1"); } @Test // DATACASS-292 @@ -260,7 +259,7 @@ public class CassandraTemplateUnitTests { assertThat(count).isEqualTo(42L); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT count(*) FROM users;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT count(1) FROM users"); } @Test // DATACASS-512 @@ -274,7 +273,7 @@ public class CassandraTemplateUnitTests { assertThat(count).isEqualTo(42L); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT COUNT(1) FROM users;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT count(1) FROM users"); } @Test // DATACASS-292, DATACASS-618 @@ -287,8 +286,8 @@ public class CassandraTemplateUnitTests { template.insert(user); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()) - .hasToString("INSERT INTO users (firstname,id,lastname) VALUES ('Walter','heisenberg','White');"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("INSERT INTO users (firstname,id,lastname) VALUES ('Walter','heisenberg','White')"); assertThat(beforeConvert).isSameAs(user); assertThat(beforeSave).isSameAs(user); } @@ -303,8 +302,8 @@ public class CassandraTemplateUnitTests { template.insert(user); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString( - "INSERT INTO vusers (firstname,id,lastname,version) VALUES ('Walter','heisenberg','White',0) IF NOT EXISTS;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo( + "INSERT INTO vusers (firstname,id,lastname,version) VALUES ('Walter','heisenberg','White',0) IF NOT EXISTS"); assertThat(beforeConvert).isSameAs(user); assertThat(beforeSave).isSameAs(user); } @@ -321,8 +320,8 @@ public class CassandraTemplateUnitTests { template.insert(user, insertOptions); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()) - .hasToString("INSERT INTO users (firstname,id,lastname) VALUES ('Walter','heisenberg','White') IF NOT EXISTS;"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("INSERT INTO users (firstname,id,lastname) VALUES ('Walter','heisenberg','White') IF NOT EXISTS"); } @Test // DATACASS-560 @@ -337,8 +336,8 @@ public class CassandraTemplateUnitTests { template.insert(user, insertOptions); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()) - .hasToString("INSERT INTO users (firstname,id,lastname) VALUES (null,'heisenberg',null);"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("INSERT INTO users (firstname,id,lastname) VALUES (NULL,'heisenberg',NULL)"); } @Test // DATACASS-292 @@ -378,8 +377,8 @@ public class CassandraTemplateUnitTests { template.update(user); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()) - .hasToString("UPDATE users SET firstname='Walter',lastname='White' WHERE id='heisenberg';"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("UPDATE users SET firstname='Walter', lastname='White' WHERE id='heisenberg'"); assertThat(beforeConvert).isSameAs(user); assertThat(beforeSave).isSameAs(user); } @@ -395,8 +394,8 @@ public class CassandraTemplateUnitTests { template.update(user); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString( - "UPDATE vusers SET firstname='Walter',lastname='White',version=1 WHERE id='heisenberg' IF version=0;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo( + "UPDATE vusers SET firstname='Walter', lastname='White', version=1 WHERE id='heisenberg' IF version=0"); assertThat(beforeConvert).isSameAs(user); assertThat(beforeSave).isSameAs(user); } @@ -414,8 +413,8 @@ public class CassandraTemplateUnitTests { assertThat(writeResult.wasApplied()).isTrue(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()) - .hasToString("UPDATE users SET firstname='Walter',lastname='White' WHERE id='heisenberg' IF EXISTS;"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("UPDATE users SET firstname='Walter', lastname='White' WHERE id='heisenberg' IF EXISTS"); } @Test // DATACASS-575 @@ -427,8 +426,8 @@ public class CassandraTemplateUnitTests { template.update(user, options); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString( - "UPDATE users SET firstname='Walter',lastname='White' WHERE id='heisenberg' IF firstname='Walter';"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("UPDATE users SET firstname='Walter', lastname='White' WHERE id='heisenberg' IF firstname='Walter'"); } @Test // DATACASS-575 @@ -440,7 +439,8 @@ public class CassandraTemplateUnitTests { template.update(query, update, User.class); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("UPDATE users SET firstname='Walter' WHERE id='heisenberg';"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("UPDATE users SET firstname='Walter' WHERE id='heisenberg'"); } @Test // DATACASS-575 @@ -456,8 +456,8 @@ public class CassandraTemplateUnitTests { template.update(query, update, User.class); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString( - "UPDATE users SET firstname='Walter' WHERE id='heisenberg' IF firstname='Walter' AND lastname='White';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo( + "UPDATE users SET firstname='Walter' WHERE id='heisenberg' IF firstname='Walter' AND lastname='White'"); } @Test // DATACASS-292 @@ -486,7 +486,7 @@ public class CassandraTemplateUnitTests { assertThat(deleted).isTrue(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("DELETE FROM users WHERE id='heisenberg';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("DELETE FROM users WHERE id='heisenberg'"); } @Test // DATACASS-292 @@ -497,7 +497,7 @@ public class CassandraTemplateUnitTests { template.delete(user); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("DELETE FROM users WHERE id='heisenberg';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("DELETE FROM users WHERE id='heisenberg'"); } @Test // DATACASS-575 @@ -509,8 +509,8 @@ public class CassandraTemplateUnitTests { template.delete(user, options); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()) - .hasToString("DELETE FROM users WHERE id='heisenberg' IF firstname='Walter';"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("DELETE FROM users WHERE id='heisenberg' IF firstname='Walter'"); } @Test // DATACASS-575 @@ -522,8 +522,8 @@ public class CassandraTemplateUnitTests { template.delete(query, User.class); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()) - .hasToString("DELETE FROM users WHERE id='heisenberg' IF firstname='Walter';"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("DELETE FROM users WHERE id='heisenberg' IF firstname='Walter'"); } @Test // DATACASS-292 @@ -547,16 +547,7 @@ public class CassandraTemplateUnitTests { template.truncate(User.class); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("TRUNCATE users;"); - } - - @Test // DATACASS-292 - @Ignore - public void batchOperationsShouldCallSession() { - - template.batchOps().insert(new User()).execute(); - - verifyNoInteractions(session); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("TRUNCATE users"); } interface UserProjection { diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/EntityQueryUtilsUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/EntityQueryUtilsUnitTests.java index a6e237d69..c2726e3e7 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/EntityQueryUtilsUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/EntityQueryUtilsUnitTests.java @@ -19,8 +19,6 @@ import static org.assertj.core.api.Assertions.*; import org.junit.Test; -import org.springframework.data.cassandra.core.convert.MappingCassandraConverter; - import com.datastax.oss.driver.api.core.CqlIdentifier; import com.datastax.oss.driver.api.core.cql.SimpleStatement; import com.datastax.oss.driver.api.querybuilder.QueryBuilder; @@ -33,35 +31,33 @@ import com.datastax.oss.driver.api.querybuilder.select.Select; */ public class EntityQueryUtilsUnitTests { - private final MappingCassandraConverter converter = new MappingCassandraConverter(); - @Test // DATACASS-106 public void shouldRetrieveTableNameFromSelect() { - Select select = QueryBuilder.selectFrom("keyspace", "table").all().where(); + Select select = QueryBuilder.selectFrom("ks", "tbl").all().where(); CqlIdentifier tableName = EntityQueryUtils.getTableName(select.build()); - assertThat(tableName).isEqualTo(CqlIdentifier.fromCql("table")); + assertThat(tableName).isEqualTo(CqlIdentifier.fromInternal("tbl")); } @Test // DATACASS-642 public void shouldRetrieveQuotedTableNameFromSelect() { - Select select = QueryBuilder.selectFrom("keyspace", "\"table\"").all().where(); + Select select = QueryBuilder.selectFrom(CqlIdentifier.fromCql("\"table\"")).all().where(); CqlIdentifier tableName = EntityQueryUtils.getTableName(select.build()); - assertThat(tableName).isEqualTo(CqlIdentifier.fromCql("table")); + assertThat(tableName).isEqualTo(CqlIdentifier.fromInternal("table")); } @Test // DATACASS-106 public void shouldRetrieveTableNameFromSimpleStatement() { assertThat(EntityQueryUtils.getTableName(SimpleStatement.newInstance("SELECT * FROM table"))) - .isEqualTo(CqlIdentifier.fromCql("table")); + .isEqualTo(CqlIdentifier.fromInternal("table")); assertThat(EntityQueryUtils.getTableName(SimpleStatement.newInstance("SELECT * FROM foo.table where"))) - .isEqualTo(CqlIdentifier.fromCql("table")); + .isEqualTo(CqlIdentifier.fromInternal("table")); } @Test // DATACASS-106 @@ -69,6 +65,6 @@ public class EntityQueryUtilsUnitTests { CqlIdentifier tableName = EntityQueryUtils.getTableName(SimpleStatement.newInstance("SELECT * from \"table\"")); - assertThat(tableName).isEqualTo(CqlIdentifier.fromCql("table")); + assertThat(tableName).isEqualTo(CqlIdentifier.fromInternal("table")); } } 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 ad054f35e..fc9ecb37a 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 @@ -43,10 +43,12 @@ import org.springframework.data.cassandra.domain.User; import org.springframework.data.cassandra.domain.VersionedUser; 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.cql.ColumnDefinition; import com.datastax.oss.driver.api.core.cql.ColumnDefinitions; 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.type.DataTypes; @@ -64,7 +66,7 @@ public class ReactiveCassandraTemplateUnitTests { @Mock ColumnDefinition columnDefinition; @Mock ColumnDefinitions columnDefinitions; - @Captor ArgumentCaptor> statementCaptor; + @Captor ArgumentCaptor statementCaptor; ReactiveCassandraTemplate template; @@ -103,7 +105,7 @@ public class ReactiveCassandraTemplateUnitTests { public void selectUsingCqlShouldReturnMappedResults() { when(reactiveResultSet.rows()).thenReturn(Flux.just(row)); - when(columnDefinitions.contains(anyString())).thenReturn(true); + when(columnDefinitions.contains(any(CqlIdentifier.class))).thenReturn(true); when(columnDefinitions.get(anyInt())).thenReturn(columnDefinition); when(columnDefinitions.firstIndexOf("id")).thenReturn(0); @@ -121,7 +123,7 @@ public class ReactiveCassandraTemplateUnitTests { .verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users"); } @Test // DATACASS-335 @@ -139,7 +141,7 @@ public class ReactiveCassandraTemplateUnitTests { public void selectOneByIdShouldReturnMappedResults() { when(reactiveResultSet.rows()).thenReturn(Flux.just(row)); - when(columnDefinitions.contains(anyString())).thenReturn(true); + when(columnDefinitions.contains(any(CqlIdentifier.class))).thenReturn(true); when(columnDefinitions.get(anyInt())).thenReturn(columnDefinition); when(columnDefinitions.firstIndexOf("id")).thenReturn(0); when(columnDefinitions.firstIndexOf("firstname")).thenReturn(1); @@ -156,14 +158,14 @@ public class ReactiveCassandraTemplateUnitTests { .verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users WHERE id='myid';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users WHERE id='myid' LIMIT 1"); } @Test // DATACASS-313 public void selectProjectedOneShouldReturnMappedResults() { when(reactiveResultSet.rows()).thenReturn(Flux.just(row)); - when(columnDefinitions.contains(anyString())).thenReturn(true); + when(columnDefinitions.contains(any(CqlIdentifier.class))).thenReturn(true); when(columnDefinitions.get(anyInt())).thenReturn(columnDefinition); when(columnDefinitions.firstIndexOf("id")).thenReturn(0); @@ -179,7 +181,7 @@ public class ReactiveCassandraTemplateUnitTests { }).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT firstname FROM users LIMIT 1;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT firstname FROM users LIMIT 1"); } @Test // DATACASS-696 @@ -187,7 +189,7 @@ public class ReactiveCassandraTemplateUnitTests { when(reactiveResultSet.rows()).thenReturn(Flux.just(row)); - template.selectOne("SELECT id FROM users WHERE id='myid';", String.class).as(StepVerifier::create) // + template.selectOne("SELECT id FROM users WHERE id='myid'", String.class).as(StepVerifier::create) // .verifyComplete(); } @@ -199,7 +201,7 @@ public class ReactiveCassandraTemplateUnitTests { template.exists("myid", User.class).as(StepVerifier::create).expectNext(true).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users WHERE id='myid';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users WHERE id='myid' LIMIT 1"); } @Test // DATACASS-335 @@ -210,7 +212,7 @@ public class ReactiveCassandraTemplateUnitTests { template.exists("myid", User.class).as(StepVerifier::create).expectNext(false).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users WHERE id='myid';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users WHERE id='myid' LIMIT 1"); } @Test // DATACASS-512 @@ -221,7 +223,7 @@ public class ReactiveCassandraTemplateUnitTests { template.exists(Query.empty(), User.class).as(StepVerifier::create).expectNext(true).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users LIMIT 1;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users LIMIT 1"); } @Test // DATACASS-512 @@ -232,7 +234,7 @@ public class ReactiveCassandraTemplateUnitTests { template.exists(Query.empty(), User.class).as(StepVerifier::create).expectNext(false).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT * FROM users LIMIT 1;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT * FROM users LIMIT 1"); } @Test // DATACASS-335 @@ -245,7 +247,7 @@ public class ReactiveCassandraTemplateUnitTests { template.count(User.class).as(StepVerifier::create).expectNext(42L).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT count(*) FROM users;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT count(1) FROM users"); } @Test // DATACASS-512 @@ -258,7 +260,7 @@ public class ReactiveCassandraTemplateUnitTests { template.count(Query.empty(), User.class).as(StepVerifier::create).expectNext(42L).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("SELECT COUNT(1) FROM users;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("SELECT count(1) FROM users"); } @Test // DATACASS-335, DATACASS-618 @@ -271,8 +273,8 @@ public class ReactiveCassandraTemplateUnitTests { template.insert(user).as(StepVerifier::create).expectNext(user).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()) - .hasToString("INSERT INTO users (firstname,id,lastname) VALUES ('Walter','heisenberg','White');"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("INSERT INTO users (firstname,id,lastname) VALUES ('Walter','heisenberg','White')"); assertThat(beforeConvert).isSameAs(user); assertThat(beforeSave).isSameAs(user); } @@ -287,8 +289,8 @@ public class ReactiveCassandraTemplateUnitTests { StepVerifier.create(template.insert(user)).expectNext(user).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString( - "INSERT INTO vusers (firstname,id,lastname,version) VALUES ('Walter','heisenberg','White',0) IF NOT EXISTS;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo( + "INSERT INTO vusers (firstname,id,lastname,version) VALUES ('Walter','heisenberg','White',0) IF NOT EXISTS"); assertThat(beforeConvert).isSameAs(user); assertThat(beforeSave).isSameAs(user); } @@ -317,8 +319,8 @@ public class ReactiveCassandraTemplateUnitTests { template.update(user).as(StepVerifier::create).expectNext(user).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()) - .hasToString("UPDATE users SET firstname='Walter',lastname='White' WHERE id='heisenberg';"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("UPDATE users SET firstname='Walter', lastname='White' WHERE id='heisenberg'"); assertThat(beforeConvert).isSameAs(user); assertThat(beforeSave).isSameAs(user); } @@ -335,8 +337,8 @@ public class ReactiveCassandraTemplateUnitTests { StepVerifier.create(template.update(user)).expectNext(user).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString( - "UPDATE vusers SET firstname='Walter',lastname='White',version=1 WHERE id='heisenberg' IF version=0;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo( + "UPDATE vusers SET firstname='Walter', lastname='White', version=1 WHERE id='heisenberg' IF version=0"); assertThat(beforeConvert).isSameAs(user); assertThat(beforeSave).isSameAs(user); } @@ -355,8 +357,8 @@ public class ReactiveCassandraTemplateUnitTests { .verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()) - .hasToString("UPDATE users SET firstname='Walter',lastname='White' WHERE id='heisenberg' IF EXISTS;"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("UPDATE users SET firstname='Walter', lastname='White' WHERE id='heisenberg' IF EXISTS"); } @Test // DATACASS-575 @@ -373,8 +375,8 @@ public class ReactiveCassandraTemplateUnitTests { .verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString( - "UPDATE users SET firstname='Walter',lastname='White' WHERE id='heisenberg' IF firstname='Walter';"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("UPDATE users SET firstname='Walter', lastname='White' WHERE id='heisenberg' IF firstname='Walter'"); } @Test // DATACASS-575 @@ -391,7 +393,8 @@ public class ReactiveCassandraTemplateUnitTests { .verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("UPDATE users SET firstname='Walter' WHERE id='heisenberg';"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("UPDATE users SET firstname='Walter' WHERE id='heisenberg'"); } @Test // DATACASS-575 @@ -412,8 +415,8 @@ public class ReactiveCassandraTemplateUnitTests { .verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString( - "UPDATE users SET firstname='Walter' WHERE id='heisenberg' IF firstname='Walter' AND lastname='White';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo( + "UPDATE users SET firstname='Walter' WHERE id='heisenberg' IF firstname='Walter' AND lastname='White'"); } @Test // DATACASS-335 @@ -427,7 +430,7 @@ public class ReactiveCassandraTemplateUnitTests { template.delete(user).as(StepVerifier::create).expectNext(user).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("DELETE FROM users WHERE id='heisenberg';"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("DELETE FROM users WHERE id='heisenberg'"); } @Test // DATACASS-575 @@ -444,8 +447,8 @@ public class ReactiveCassandraTemplateUnitTests { .verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()) - .hasToString("DELETE FROM users WHERE id='heisenberg' IF firstname='Walter';"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("DELETE FROM users WHERE id='heisenberg' IF firstname='Walter'"); } @Test // DATACASS-575 @@ -462,8 +465,8 @@ public class ReactiveCassandraTemplateUnitTests { .verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()) - .hasToString("DELETE FROM users WHERE id='heisenberg' IF firstname='Walter';"); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("DELETE FROM users WHERE id='heisenberg' IF firstname='Walter'"); } @Test // DATACASS-335 @@ -472,7 +475,7 @@ public class ReactiveCassandraTemplateUnitTests { template.truncate(User.class).as(StepVerifier::create).verifyComplete(); verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue()).hasToString("TRUNCATE users;"); + assertThat(statementCaptor.getValue().getQuery()).isEqualTo("TRUNCATE users"); } interface UserProjection { diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/StatementFactoryUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/StatementFactoryUnitTests.java index c1b7a147c..ee1350193 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/StatementFactoryUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/StatementFactoryUnitTests.java @@ -17,17 +17,22 @@ package org.springframework.data.cassandra.core; import static org.assertj.core.api.Assertions.*; +import java.time.Duration; +import java.util.Collections; import java.util.List; import java.util.Map; import java.util.Set; +import org.junit.Ignore; import org.junit.Test; import org.springframework.data.annotation.Id; import org.springframework.data.cassandra.core.convert.CassandraConverter; import org.springframework.data.cassandra.core.convert.MappingCassandraConverter; import org.springframework.data.cassandra.core.convert.UpdateMapper; +import org.springframework.data.cassandra.core.cql.WriteOptions; import org.springframework.data.cassandra.core.cql.util.StatementBuilder; +import org.springframework.data.cassandra.core.cql.util.StatementBuilder.ParameterHandling; import org.springframework.data.cassandra.core.mapping.CassandraPersistentEntity; import org.springframework.data.cassandra.core.mapping.Column; import org.springframework.data.cassandra.core.query.Columns; @@ -38,6 +43,7 @@ import org.springframework.data.cassandra.domain.Group; import org.springframework.data.domain.Sort; import com.datastax.oss.driver.api.querybuilder.delete.Delete; +import com.datastax.oss.driver.api.querybuilder.insert.RegularInsert; import com.datastax.oss.driver.api.querybuilder.select.Select; /** @@ -62,7 +68,7 @@ public class StatementFactoryUnitTests { StatementBuilder select = statementFactory.select(query, groupEntity); - assertThat(select.build().toString()).isEqualTo("SELECT age FROM group WHERE foo='bar';"); + assertThat(select.build(ParameterHandling.INLINE).getQuery()).isEqualTo("SELECT age FROM group WHERE foo='bar'"); } @Test // DATACASS-549 @@ -82,7 +88,7 @@ public class StatementFactoryUnitTests { StatementBuilder select = statementFactory.select(query, groupEntity); - assertThat(select.build().toString()).isEqualTo("SELECT age FROM group WHERE foo IS NOT NULL;"); + assertThat(select.build(ParameterHandling.INLINE).getQuery()) + .isEqualTo("SELECT age FROM group WHERE foo IS NOT NULL"); } @Test // DATACASS-343 @@ -103,7 +110,7 @@ public class StatementFactoryUnitTests { StatementBuilder select = statementFactory.select(query, converter.getMappingContext().getRequiredPersistentEntity(Group.class)); - assertThat(select.build().toString()) - .isEqualTo("SELECT * FROM group ORDER BY hash_prefix ASC LIMIT 10 ALLOW FILTERING;"); + assertThat(select.build(ParameterHandling.INLINE).getQuery()) + .isEqualTo("SELECT * FROM group ORDER BY hash_prefix ASC LIMIT 10 ALLOW FILTERING"); } @Test // DATACASS-343 @@ -126,18 +133,83 @@ public class StatementFactoryUnitTests { StatementBuilder delete = statementFactory.delete(query, converter.getMappingContext().getRequiredPersistentEntity(Group.class)); - assertThat(delete.build().toString()).isEqualTo("DELETE age FROM group;"); + assertThat(delete.build(ParameterHandling.INLINE).getQuery()).isEqualTo("DELETE age FROM group"); } @Test // DATACASS-343 - public void shouldMapDeleteQueryWithTtlColumns() { + public void shouldMapDeleteQueryWithTimestampColumns() { - Query query = Query.query(Criteria.where("foo").is("bar")); + DeleteOptions options = DeleteOptions.builder().timestamp(1234).build(); + Query query = Query.query(Criteria.where("foo").is("bar")).queryOptions(options); StatementBuilder delete = statementFactory.delete(query, converter.getMappingContext().getRequiredPersistentEntity(Group.class)); - assertThat(delete.build().toString()).isEqualTo("DELETE FROM group WHERE foo='bar';"); + assertThat(delete.build(ParameterHandling.INLINE).getQuery()) + .isEqualTo("DELETE FROM group USING TIMESTAMP 1234 WHERE foo='bar'"); + } + + @Test // DATACASS-656 + public void shouldCreateInsert() { + + Person person = new Person(); + person.id = "foo"; + + StatementBuilder insert = statementFactory.insert(person, WriteOptions.empty()); + + assertThat(insert.build(ParameterHandling.INLINE).getQuery()).isEqualTo("INSERT INTO person (id) VALUES ('foo')"); + } + + @Test // DATACASS-656 + public void shouldCreateInsertIfNotExists() { + + InsertOptions options = InsertOptions.builder().withIfNotExists().build(); + Person person = new Person(); + person.id = "foo"; + + StatementBuilder insert = statementFactory.insert(person, options); + + assertThat(insert.build(ParameterHandling.INLINE).getQuery()) + .isEqualTo("INSERT INTO person (id) VALUES ('foo') IF NOT EXISTS"); + } + + @Test // DATACASS-656 + public void shouldCreateSetInsertNulls() { + + InsertOptions options = InsertOptions.builder().withInsertNulls().build(); + Person person = new Person(); + person.id = "foo"; + + StatementBuilder insert = statementFactory.insert(person, options); + + assertThat(insert.build(ParameterHandling.INLINE).getQuery()).isEqualTo( + "INSERT INTO person (first_name,id,list,map,number,set_col) VALUES (NULL,'foo',NULL,NULL,NULL,NULL)"); + } + + @Test // DATACASS-656 + public void shouldCreateSetInsertWithTtl() { + + WriteOptions options = WriteOptions.builder().ttl(Duration.ofMinutes(1)).build(); + Person person = new Person(); + person.id = "foo"; + + StatementBuilder insert = statementFactory.insert(person, options); + + assertThat(insert.build(ParameterHandling.INLINE).getQuery()) + .isEqualTo("INSERT INTO person (id) VALUES ('foo') USING TTL 60"); + } + + @Test // DATACASS-656 + public void shouldCreateSetInsertWithTimestamp() { + + WriteOptions options = WriteOptions.builder().timestamp(1234).build(); + Person person = new Person(); + person.id = "foo"; + + StatementBuilder insert = statementFactory.insert(person, options); + + assertThat(insert.build(ParameterHandling.INLINE).getQuery()) + .isEqualTo("INSERT INTO person (id) VALUES ('foo') USING TIMESTAMP 1234"); } @Test // DATACASS-343 @@ -148,16 +220,44 @@ public class StatementFactoryUnitTests { StatementBuilder update = statementFactory.update(query, Update.empty().set("firstName", "baz").set("boo", "baa"), personEntity); - assertThat(update.build().toString()).isEqualTo("UPDATE person SET first_name='baz',boo='baa' WHERE foo='bar';"); + assertThat(update.build(ParameterHandling.INLINE).getQuery()) + .isEqualTo("UPDATE person SET first_name='baz', boo='baa' WHERE foo='bar'"); + } + + @Test // DATACASS-656 + public void shouldCreateSetUpdateWithTtl() { + + WriteOptions options = WriteOptions.builder().ttl(Duration.ofMinutes(1)).build(); + Query query = Query.query(Criteria.where("foo").is("bar")).queryOptions(options); + + StatementBuilder update = statementFactory.update(query, + Update.empty().set("firstName", "baz"), personEntity); + + assertThat(update.build(ParameterHandling.INLINE).getQuery()) + .isEqualTo("UPDATE person USING TTL 60 SET first_name='baz' WHERE foo='bar'"); + } + + @Test // DATACASS-656 + public void shouldCreateSetUpdateWithTimestamp() { + + WriteOptions options = WriteOptions.builder().timestamp(1234).build(); + Query query = Query.query(Criteria.where("foo").is("bar")).queryOptions(options); + + StatementBuilder update = statementFactory.update(query, + Update.empty().set("firstName", "baz"), personEntity); + + assertThat(update.build(ParameterHandling.INLINE).getQuery()) + .isEqualTo("UPDATE person USING TIMESTAMP 1234 SET first_name='baz' WHERE foo='bar'"); } @Test // DATACASS-343 + @Ignore("No operator for set at index yet") public void shouldCreateSetAtIndexUpdate() { StatementBuilder update = statementFactory .update(Query.empty(), Update.empty().set("list").atIndex(10).to("Euro"), personEntity); - assertThat(update.build().toString()).isEqualTo("UPDATE person SET list[10]='Euro';"); + assertThat(update.build(ParameterHandling.INLINE).getQuery()).isEqualTo("UPDATE person SET list[10]='Euro'"); } @Test // DATACASS-343 @@ -166,7 +266,7 @@ public class StatementFactoryUnitTests { StatementBuilder update = statementFactory .update(Query.empty(), Update.empty().set("map").atKey("baz").to("Euro"), personEntity); - assertThat(update.build().toString()).isEqualTo("UPDATE person SET map['baz']='Euro';"); + assertThat(update.build(ParameterHandling.INLINE).getQuery()).isEqualTo("UPDATE person SET map['baz']='Euro'"); } @Test // DATACASS-343 @@ -175,25 +275,28 @@ public class StatementFactoryUnitTests { StatementBuilder update = statementFactory .update(Query.empty(), Update.empty().addTo("map").entry("foo", "Euro"), personEntity); - assertThat(update.build().toString()).isEqualTo("UPDATE person SET map=map+{'foo':'Euro'};"); + assertThat(update.build(ParameterHandling.INLINE).getQuery()).isEqualTo("UPDATE person SET map+={'foo':'Euro'}"); } @Test // DATACASS-343 + @Ignore("Missing operator") public void shouldPrependAllToList() { StatementBuilder update = statementFactory .update(Query.empty(), Update.empty().addTo("list").prependAll("foo", "Euro"), personEntity); - assertThat(update.build().toString()).isEqualTo("UPDATE person SET list=['foo','Euro']+list;"); + assertThat(update.build(ParameterHandling.INLINE).getQuery()) + .isEqualTo("UPDATE person SET list=['foo','Euro']+list"); } @Test // DATACASS-343 + @Ignore("Missing operator") public void shouldAppendAllToList() { StatementBuilder update = statementFactory .update(Query.empty(), Update.empty().addTo("list").appendAll("foo", "Euro"), personEntity); - assertThat(update.toString()).isEqualTo("UPDATE person SET list=list+['foo','Euro'];"); + assertThat(update.build(ParameterHandling.INLINE).getQuery()).isEqualTo("UPDATE person SET list+=['foo','Euro']"); } @Test // DATACASS-343 @@ -202,7 +305,7 @@ public class StatementFactoryUnitTests { StatementBuilder update = statementFactory .update(Query.empty(), Update.empty().remove("list", "Euro"), personEntity); - assertThat(update.build().toString()).isEqualTo("UPDATE person SET list=list-['Euro'];"); + assertThat(update.build(ParameterHandling.INLINE).getQuery()).isEqualTo("UPDATE person SET list-=['Euro']"); } @Test // DATACASS-343 @@ -211,16 +314,18 @@ public class StatementFactoryUnitTests { StatementBuilder update = statementFactory .update(Query.empty(), Update.empty().clear("list"), personEntity); - assertThat(update.build().toString()).isEqualTo("UPDATE person SET list=[];"); + assertThat(update.build(ParameterHandling.INLINE).getQuery()).isEqualTo("UPDATE person SET list=[]"); } @Test // DATACASS-343 + @Ignore("Missing operator") public void shouldAddAllToSet() { StatementBuilder update = statementFactory .update(Query.empty(), Update.empty().addTo("set").appendAll("foo", "Euro"), personEntity); - assertThat(update.build().toString()).isEqualTo("UPDATE person SET set_col=set_col+{'foo','Euro'};"); + assertThat(update.build(ParameterHandling.INLINE).getQuery()) + .isEqualTo("UPDATE person SET set_col+={'foo','Euro'}"); } @Test // DATACASS-343 @@ -229,7 +334,7 @@ public class StatementFactoryUnitTests { StatementBuilder update = statementFactory .update(Query.empty(), Update.empty().remove("set", "Euro"), personEntity); - assertThat(update.build().toString()).isEqualTo("UPDATE person SET set_col=set_col-{'Euro'};"); + assertThat(update.build(ParameterHandling.INLINE).getQuery()).isEqualTo("UPDATE person SET set_col-={'Euro'}"); } @Test // DATACASS-343 @@ -238,7 +343,7 @@ public class StatementFactoryUnitTests { StatementBuilder update = statementFactory .update(Query.empty(), Update.empty().clear("set"), personEntity); - assertThat(update.build().toString()).isEqualTo("UPDATE person SET set_col={};"); + assertThat(update.build(ParameterHandling.INLINE).getQuery()).isEqualTo("UPDATE person SET set_col={}"); } @Test // DATACASS-343 @@ -247,7 +352,7 @@ public class StatementFactoryUnitTests { StatementBuilder update = statementFactory .update(Query.empty(), Update.empty().increment("number"), personEntity); - assertThat(update.build().toString()).isEqualTo("UPDATE person SET number=number+1;"); + assertThat(update.build(ParameterHandling.INLINE).getQuery()).isEqualTo("UPDATE person SET number+=1"); } @Test // DATACASS-343 @@ -256,7 +361,7 @@ public class StatementFactoryUnitTests { StatementBuilder update = statementFactory .update(Query.empty(), Update.empty().decrement("number"), personEntity); - assertThat(update.build().toString()).isEqualTo("UPDATE person SET number=number-1;"); + assertThat(update.build(ParameterHandling.INLINE).getQuery()).isEqualTo("UPDATE person SET number-=1"); } @Test // DATACASS-569 @@ -268,7 +373,104 @@ public class StatementFactoryUnitTests { StatementBuilder update = statementFactory.update(query, Update.empty().set("firstName", "baz"), personEntity); - assertThat(update.build().toString()).isEqualTo("UPDATE person SET first_name='baz' WHERE foo='bar' IF EXISTS;"); + assertThat(update.build(ParameterHandling.INLINE).getQuery()) + .isEqualTo("UPDATE person SET first_name='baz' WHERE foo='bar' IF EXISTS"); + } + + @Test // DATACASS-656 + public void shouldCreateSetUpdateIfCondition() { + + Query query = Query.query(Criteria.where("foo").is("bar")) + .queryOptions(UpdateOptions.builder().ifCondition(Criteria.where("foo").is("baz")).build()); + + StatementBuilder update = statementFactory.update(query, + Update.empty().set("firstName", "baz"), personEntity); + + assertThat(update.build(ParameterHandling.INLINE).getQuery()) + .isEqualTo("UPDATE person SET first_name='baz' WHERE foo='bar' IF foo='baz'"); + } + + @Test // DATACASS-656 + public void shouldCreateSetUpdateFromObject() { + + Person person = new Person(); + person.id = "foo"; + person.firstName = "bar"; + + StatementBuilder update = statementFactory.update(person, + WriteOptions.empty()); + + assertThat(update.build(ParameterHandling.INLINE).getQuery()) + .isEqualTo("UPDATE person SET first_name='bar', list=NULL, map=NULL, number=NULL, set_col=NULL WHERE id='foo'"); + } + + @Test // DATACASS-656 + public void shouldCreateSetUpdateFromObjectIfExists() { + + UpdateOptions options = UpdateOptions.builder().withIfExists().build(); + Person person = new Person(); + person.id = "foo"; + person.firstName = "bar"; + + StatementBuilder update = statementFactory.update(person, + options); + + assertThat(update.build(ParameterHandling.INLINE).getQuery()).endsWith("IF EXISTS"); + } + + @Test // DATACASS-656 + public void shouldCreateSetUpdateFromObjectIfCondition() { + + UpdateOptions options = UpdateOptions.builder().ifCondition(Criteria.where("foo").is("bar")).build(); + Person person = new Person(); + person.id = "foo"; + person.firstName = "bar"; + + StatementBuilder update = statementFactory.update(person, + options); + + assertThat(update.build(ParameterHandling.INLINE).getQuery()).endsWith("IF foo='bar'"); + } + + @Test // DATACASS-656 + public void shouldCreateSetUpdateFromObjectWithTtl() { + + WriteOptions options = WriteOptions.builder().ttl(Duration.ofMinutes(1)).build(); + Person person = new Person(); + person.id = "foo"; + + StatementBuilder update = statementFactory.update(person, + options); + + assertThat(update.build(ParameterHandling.INLINE).getQuery()).startsWith("UPDATE person USING TTL 60 SET"); + } + + @Test // DATACASS-656 + public void shouldCreateSetUpdateFromObjectWithTimestamp() { + + WriteOptions options = WriteOptions.builder().timestamp(1234).build(); + Person person = new Person(); + person.id = "foo"; + + StatementBuilder update = statementFactory.update(person, + options); + + assertThat(update.build(ParameterHandling.INLINE).getQuery()).startsWith("UPDATE person USING TIMESTAMP 1234 SET"); + } + + @Test // DATACASS-656 + public void shouldCreateSetUpdateFromObjectWithEmptyCollections() { + + Person person = new Person(); + person.id = "foo"; + person.set = Collections.emptySet(); + person.list = Collections.emptyList(); + + StatementBuilder update = statementFactory.update(person, + WriteOptions.empty()); + + assertThat(update.build(ParameterHandling.INLINE).getQuery()) + .isEqualTo("UPDATE person SET first_name=NULL, list=[], map=NULL, number=NULL, set_col={} WHERE id='foo'"); } @Test // DATACASS-512 @@ -279,7 +481,8 @@ public class StatementFactoryUnitTests { StatementBuilder