From ab730211871559e702dcececdc526498f2257f2b Mon Sep 17 00:00:00 2001 From: John Blum Date: Fri, 1 Jul 2016 17:26:04 -0700 Subject: [PATCH] DATACASS-310 - Fix CqlTemplate and CassandraTemplate returning null for data access operations returning a Collection. --- .../cassandra/core/CqlTemplate.java | 106 ++++++----- .../core/support/EmptyResultSet.java | 121 +++++++++++++ .../core/support/ResultSetAdapter.java | 137 ++++++++++++++ .../CassandraReadTimeoutException.java | 12 +- .../cassandra/core/CqlTemplateUnitTests.java | 106 ++++++++--- .../core/support/EmptyResultSetUnitTests.java | 109 +++++++++++ .../cassandra/core/CassandraTemplate.java | 150 +++++++--------- .../core/CassandraTemplateUnitTests.java | 170 ++++++++++++++++++ 8 files changed, 755 insertions(+), 156 deletions(-) create mode 100644 spring-cql/src/main/java/org/springframework/cassandra/core/support/EmptyResultSet.java create mode 100644 spring-cql/src/main/java/org/springframework/cassandra/core/support/ResultSetAdapter.java create mode 100644 spring-cql/src/test/java/org/springframework/cassandra/core/support/EmptyResultSetUnitTests.java create mode 100644 spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateUnitTests.java diff --git a/spring-cql/src/main/java/org/springframework/cassandra/core/CqlTemplate.java b/spring-cql/src/main/java/org/springframework/cassandra/core/CqlTemplate.java index 3c5bf0cce..a7f794136 100644 --- a/spring-cql/src/main/java/org/springframework/cassandra/core/CqlTemplate.java +++ b/spring-cql/src/main/java/org/springframework/cassandra/core/CqlTemplate.java @@ -124,6 +124,34 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { return cql; } + protected T logStatement(T statement) { + logDebug("executing statement [{}]", statement); + return statement; + } + + /** + * Add common {@link QueryOptions} to Cassandra {@link PreparedStatement}s. + * + * @param preparedStatement the Cassandra {@link PreparedStatement} to execute. + * @param queryOptions query options (e.g. consistency level) to add to the Cassandra {@link PreparedStatement}. + */ + public static PreparedStatement addPreparedStatementOptions(PreparedStatement preparedStatement, + QueryOptions queryOptions) { + + if (queryOptions != null) { + + if (queryOptions.getConsistencyLevel() != null) { + preparedStatement.setConsistencyLevel(ConsistencyLevelResolver.resolve(queryOptions.getConsistencyLevel())); + } + + if (queryOptions.getRetryPolicy() != null) { + preparedStatement.setRetryPolicy(RetryPolicyResolver.resolve(queryOptions.getRetryPolicy())); + } + } + + return preparedStatement; + } + /** * Add common {@link QueryOptions} to all types of queries. * @@ -189,29 +217,6 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { return update; } - /** - * Add common {@link QueryOptions} to Cassandra {@link PreparedStatement}s. - * - * @param preparedStatement the Cassandra {@link PreparedStatement} to execute. - * @param queryOptions query options (e.g. consistency level) to add to the Cassandra {@link PreparedStatement}. - */ - public static PreparedStatement addPreparedStatementOptions(PreparedStatement preparedStatement, - QueryOptions queryOptions) { - - if (queryOptions != null) { - - if (queryOptions.getConsistencyLevel() != null) { - preparedStatement.setConsistencyLevel(ConsistencyLevelResolver.resolve(queryOptions.getConsistencyLevel())); - } - - if (queryOptions.getRetryPolicy() != null) { - preparedStatement.setRetryPolicy(RetryPolicyResolver.resolve(queryOptions.getRetryPolicy())); - } - } - - return preparedStatement; - } - /** * Constructs an uninitialized instance of {@link CqlTemplate}. A Cassandra {@link Session} is required before use. * @@ -227,7 +232,7 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { * @see com.datastax.driver.core.Session * @see #setSession(Session) */ - // TODO should probably not call setSession(..) in constructor for initialization safety; + // TODO: should not call setSession(..) in constructor for initialization safety; // only really matters if CqlTemplate makes Thread-safety guarantees, which currently it does not. public CqlTemplate(Session session) { setSession(session); @@ -242,15 +247,32 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { */ protected T doExecute(SessionCallback callback) { - Assert.notNull(callback); + Assert.notNull(callback, "SessionCallback must not be null"); try { return callback.doInSession(getSession()); - } catch (DataAccessException e) { + } catch (Exception e) { throw translateExceptionIfPossible(e); } } + protected ResultSet doExecuteQueryReturnResultSet(final String query) { + return doExecute(new SessionCallback() { + @Override + public ResultSet doInSession(Session session) throws DataAccessException { + return session.execute(logCql(query)); + } + }); + } + + protected ResultSet doExecuteQueryReturnResultSet(final Select select) { + return doExecute(new SessionCallback() { + @Override public ResultSet doInSession(Session session) throws DataAccessException { + return session.execute(logStatement(select)); + } + }); + } + @Override public T execute(SessionCallback sessionCallback) { return doExecute(sessionCallback); @@ -960,6 +982,7 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { @Override public T query(String cql, PreparedStatementBinder preparedStatementBinder, ResultSetExtractor resultSetExtractor) { + return query(cql, preparedStatementBinder, resultSetExtractor, null); } @@ -974,6 +997,7 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { @Override public void query(String cql, PreparedStatementBinder preparedStatementBinder, RowCallbackHandler rowCallbackHandler) { + query(cql, preparedStatementBinder, rowCallbackHandler, null); } @@ -992,13 +1016,15 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { @Override public List query(String cql, PreparedStatementBinder preparedStatementBinder, RowMapper rowMapper, QueryOptions queryOptions) { + return query(new CachedPreparedStatementCreator(logCql(cql)), preparedStatementBinder, rowMapper, queryOptions); } @Override public void ingest(String cql, RowIterator rowIterator, WriteOptions options) { - CachedPreparedStatementCreator cachedPreparedStatementCreator = new CachedPreparedStatementCreator(logCql(cql)); + CachedPreparedStatementCreator cachedPreparedStatementCreator = + new CachedPreparedStatementCreator(logCql(cql)); PreparedStatement preparedStatement = addPreparedStatementOptions( cachedPreparedStatementCreator.createPreparedStatement(getSession()), options); @@ -1023,8 +1049,8 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { @Override public void ingest(String cql, final List> rows, WriteOptions writeOptions) { - Assert.notNull(rows); - Assert.notEmpty(rows); + Assert.notNull(rows, "Rows must not be null"); + Assert.notEmpty(rows, "Rows must not be empty"); ingest(cql, new RowIterator() { @@ -1056,17 +1082,17 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { int index = 0; @Override - public Object[] next() { - - if (!hasNext()) { - throw new NoSuchElementException("No more elements"); - } - return rows[index++]; + public boolean hasNext() { + return (index < rows.length); } @Override - public boolean hasNext() { - return (index < rows.length); + public Object[] next() { + if (!hasNext()) { + throw new NoSuchElementException("No more elements"); + } + + return rows[index++]; } }, writeOptions); } @@ -1160,6 +1186,7 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { @Override public List query(PreparedStatementCreator preparedStatementCreator, PreparedStatementBinder preparedStatementBinder, RowMapper rowMapper) { + return query(preparedStatementCreator, preparedStatementBinder, rowMapper, null); } @@ -1350,19 +1377,16 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { @Override public Cancellable executeAsynchronously(Insert insert, AsynchronousQueryListener listener) { - return doExecuteAsync(insert, listener); } @Override public Cancellable executeAsynchronously(Truncate truncate, AsynchronousQueryListener listener) { - return doExecuteAsync(truncate, listener); } @Override public Cancellable executeAsynchronously(Update update, AsynchronousQueryListener listener) { - return doExecuteAsync(update, listener); } @@ -1572,7 +1596,6 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { @Override public Cancellable queryForMapAsynchronously(String cql, QueryForMapListener listener) { - return queryForMapAsynchronously(cql, listener, null); } @@ -1709,7 +1732,6 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { @Override public ResultSet getResultSetUninterruptibly(ResultSetFuture resultSetFuture, long timeout, TimeUnit timeUnit) { - try { timeUnit = (timeUnit != null ? timeUnit : TimeUnit.MILLISECONDS); diff --git a/spring-cql/src/main/java/org/springframework/cassandra/core/support/EmptyResultSet.java b/spring-cql/src/main/java/org/springframework/cassandra/core/support/EmptyResultSet.java new file mode 100644 index 000000000..180d56c73 --- /dev/null +++ b/spring-cql/src/main/java/org/springframework/cassandra/core/support/EmptyResultSet.java @@ -0,0 +1,121 @@ +/* + * Copyright 2013-2016 the original author or authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cassandra.core.support; + +import java.util.Collections; +import java.util.Iterator; +import java.util.List; + +import com.datastax.driver.core.ExecutionInfo; +import com.datastax.driver.core.ResultSet; +import com.datastax.driver.core.Row; + +/** + * An empty {@link com.datastax.driver.core.ResultSet} implementation + * + * @author John Blum + * @see org.springframework.cassandra.core.support.ResultSetAdapter + * @see com.datastax.driver.core.ResultSet + * @since 1.5.0 + */ +public class EmptyResultSet extends ResultSetAdapter { + + protected static final EmptyResultSet INSTANCE = new EmptyResultSet(); + + /** + * Returns the given {@link ResultSet} if not null, otherwise returns an empty {@link ResultSet}. + * + * @param resultSet {@link ResultSet} to evaluate for {@literal null}. + * @return the given {@link ResultSet} if not null, otherwise return an empty {@link ResultSet}. + * @see com.datastax.driver.core.ResultSet + */ + public static ResultSet nullSafeResultSet(ResultSet resultSet) { + return (resultSet != null ? resultSet : EmptyResultSet.INSTANCE); + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#isExhausted() + */ + @Override + public boolean isExhausted() { + return true; + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#isFullyFetched() + */ + @Override + public boolean isFullyFetched() { + return true; + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#getAvailableWithoutFetching() + */ + @Override + public int getAvailableWithoutFetching() { + return 0; + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#getAllExecutionInfo() + */ + @Override + public List getAllExecutionInfo() { + return Collections.emptyList(); + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#getExecutionInfo() + */ + @Override + public ExecutionInfo getExecutionInfo() { + return null; + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#all() + */ + @Override + public List all() { + return Collections.emptyList(); + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#iterator() + */ + @Override + public Iterator iterator() { + return Collections.emptyIterator(); + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#one() + */ + @Override + public Row one() { + return null; + } +} diff --git a/spring-cql/src/main/java/org/springframework/cassandra/core/support/ResultSetAdapter.java b/spring-cql/src/main/java/org/springframework/cassandra/core/support/ResultSetAdapter.java new file mode 100644 index 000000000..62c94a19b --- /dev/null +++ b/spring-cql/src/main/java/org/springframework/cassandra/core/support/ResultSetAdapter.java @@ -0,0 +1,137 @@ +/* + * Copyright 2013-2016 the original author or authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cassandra.core.support; + +import java.util.Iterator; +import java.util.List; + +import com.datastax.driver.core.ColumnDefinitions; +import com.datastax.driver.core.ExecutionInfo; +import com.datastax.driver.core.ResultSet; +import com.datastax.driver.core.Row; +import com.google.common.util.concurrent.ListenableFuture; + +/** + * An Adapter class to simply implementations of the {@link ResultSet} interface. + * + * @author John Blum + * @see com.datastax.driver.core.ResultSet + * @since 1.5.0 + */ +public class ResultSetAdapter implements ResultSet { + + private static final String NOT_SUPPORTED = "Not Supported"; + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#isExhausted() + */ + @Override + public boolean isExhausted() { + throw new UnsupportedOperationException(NOT_SUPPORTED); + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#isFullyFetched() + */ + @Override + public boolean isFullyFetched() { + throw new UnsupportedOperationException(NOT_SUPPORTED); + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#getAvailableWithoutFetching() + */ + @Override + public int getAvailableWithoutFetching() { + throw new UnsupportedOperationException(NOT_SUPPORTED); + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#getColumnDefinitions() + */ + @Override + public ColumnDefinitions getColumnDefinitions() { + throw new UnsupportedOperationException(NOT_SUPPORTED); + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#getAllExecutionInfo() + */ + @Override + public List getAllExecutionInfo() { + throw new UnsupportedOperationException(NOT_SUPPORTED); + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#getExecutionInfo() + */ + @Override + public ExecutionInfo getExecutionInfo() { + throw new UnsupportedOperationException(NOT_SUPPORTED); + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#all() + */ + @Override + public List all() { + throw new UnsupportedOperationException(NOT_SUPPORTED); + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#fetchMoreResults() + */ + @Override + public ListenableFuture fetchMoreResults() { + throw new UnsupportedOperationException(NOT_SUPPORTED); + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#iterator() + */ + @Override + public Iterator iterator() { + throw new UnsupportedOperationException(NOT_SUPPORTED); + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#one() + */ + @Override + public Row one() { + throw new UnsupportedOperationException(NOT_SUPPORTED); + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.ResultSet#wasApplied() + */ + @Override + public boolean wasApplied() { + throw new UnsupportedOperationException(NOT_SUPPORTED); + } +} diff --git a/spring-cql/src/main/java/org/springframework/cassandra/support/exception/CassandraReadTimeoutException.java b/spring-cql/src/main/java/org/springframework/cassandra/support/exception/CassandraReadTimeoutException.java index 095cd5fcc..9faca987e 100644 --- a/spring-cql/src/main/java/org/springframework/cassandra/support/exception/CassandraReadTimeoutException.java +++ b/spring-cql/src/main/java/org/springframework/cassandra/support/exception/CassandraReadTimeoutException.java @@ -1,12 +1,12 @@ /* * Copyright 2013-2014 the original author or authors. - * + * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. * You may obtain a copy of the License at - * + * * http://www.apache.org/licenses/LICENSE-2.0 - * + * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. @@ -19,7 +19,7 @@ import org.springframework.dao.QueryTimeoutException; /** * Spring data access exception for a Cassandra read timeout. - * + * * @author Matthew T. Adams */ public class CassandraReadTimeoutException extends QueryTimeoutException { @@ -28,8 +28,8 @@ public class CassandraReadTimeoutException extends QueryTimeoutException { private boolean wasDataReceived; - public CassandraReadTimeoutException(boolean wasDataReceived, String msg, Throwable cause) { - super(msg); + public CassandraReadTimeoutException(boolean wasDataReceived, String message, Throwable cause) { + super(message, cause); this.wasDataReceived = wasDataReceived; } diff --git a/spring-cql/src/test/java/org/springframework/cassandra/core/CqlTemplateUnitTests.java b/spring-cql/src/test/java/org/springframework/cassandra/core/CqlTemplateUnitTests.java index 4ecb69875..c58b045ab 100644 --- a/spring-cql/src/test/java/org/springframework/cassandra/core/CqlTemplateUnitTests.java +++ b/spring-cql/src/test/java/org/springframework/cassandra/core/CqlTemplateUnitTests.java @@ -28,12 +28,18 @@ import org.junit.rules.ExpectedException; import org.junit.runner.RunWith; import org.mockito.Mock; import org.mockito.runners.MockitoJUnitRunner; +import org.springframework.cassandra.support.CassandraExceptionTranslator; +import org.springframework.cassandra.support.exception.CassandraReadTimeoutException; +import org.springframework.dao.DataAccessException; import org.springframework.dao.IncorrectResultSizeDataAccessException; import com.datastax.driver.core.ColumnDefinitions; +import com.datastax.driver.core.ConsistencyLevel; import com.datastax.driver.core.ResultSet; import com.datastax.driver.core.Row; import com.datastax.driver.core.Session; +import com.datastax.driver.core.exceptions.ReadTimeoutException; +import com.datastax.driver.core.querybuilder.Select; /** * The CqlTemplateUnitTests class is a test suite of test cases testing the contract and functionality of the @@ -46,19 +52,86 @@ import com.datastax.driver.core.Session; @SuppressWarnings("unchecked") public class CqlTemplateUnitTests { - @Rule public ExpectedException exception = ExpectedException.none(); + @Rule + public ExpectedException exception = ExpectedException.none(); private CqlTemplate template; - @Mock private Session mockSession; + @Mock + private Session mockSession; @Before public void setup() { template = new CqlTemplate(mockSession); + template.setExceptionTranslator(new CassandraExceptionTranslator()); + } + + @Test + public void doExecuteInSessionCallbackIsCalled() { + String result = template.doExecute(new SessionCallback() { + @Override public String doInSession(Session session) throws DataAccessException { + session.execute("test"); + return "test"; + } + }); + + assertThat(result, is(equalTo("test"))); + + verify(mockSession, times(1)).execute(eq("test")); + } + + @Test + public void doExecuteInSessionCallbackTranslatesException() { + exception.expect(CassandraReadTimeoutException.class); + exception.expectCause(org.hamcrest.Matchers.isA(ReadTimeoutException.class)); + + template.doExecute(new SessionCallback() { + @Override public String doInSession(Session session) throws DataAccessException { + throw new ReadTimeoutException(ConsistencyLevel.ALL, 0, 1, true); + } + }); + } + + @Test + public void doExecuteWithNullSessionCallbackThrowsIllegalArgumentException() { + exception.expect(IllegalArgumentException.class); + exception.expectCause(is(nullValue(Throwable.class))); + exception.expectMessage("SessionCallback must not be null"); + + template.doExecute((SessionCallback) null); + } + + @Test + public void doExecuteQueryReturnsResultSetForOqlQueryString() { + ResultSet mockResultSet = mock(ResultSet.class); + + when(mockSession.execute(eq("SELECT * FROM Customers"))).thenReturn(mockResultSet); + + ResultSet resultSet = template.doExecuteQueryReturnResultSet("SELECT * FROM Customers"); + + assertThat(resultSet, is(equalTo(mockResultSet))); + + verify(mockSession, times(1)).execute(eq("SELECT * FROM Customers")); + verifyZeroInteractions(mockResultSet); + } + + @Test + public void doExecuteSelectReturnsResultSetForOqlQueryString() { + Select mockSelect = mock(Select.class); + ResultSet mockResultSet = mock(ResultSet.class); + + when(mockSession.execute(eq(mockSelect))).thenReturn(mockResultSet); + + ResultSet resultSet = template.doExecuteQueryReturnResultSet(mockSelect); + + assertThat(resultSet, is(equalTo(mockResultSet))); + + verify(mockSession, times(1)).execute(eq(mockSelect)); + verifyZeroInteractions(mockResultSet); } /** - * @see DATACASS-286 + * @see DATACASS-286 */ @Test public void firstColumnToObjectReturnsColumnValue() { @@ -74,12 +147,11 @@ public class CqlTemplateUnitTests { when(mockIterator.next()).thenReturn(mockColumnDefinition); template = new CqlTemplate() { - @Override T columnToObject(Row row, ColumnDefinitions.Definition columnDefinition) { - assertThat(row, is(sameInstance(mockRow))); assertThat(columnDefinition, is(sameInstance(mockColumnDefinition))); + return (T) "test"; } }; @@ -94,7 +166,7 @@ public class CqlTemplateUnitTests { } /** - * @see DATACASS-286 + * @see DATACASS-286 */ @Test public void firstColumnToObjectReturnsNull() { @@ -116,7 +188,7 @@ public class CqlTemplateUnitTests { } /** - * @see DATACASS-286 + * @see DATACASS-286 */ @Test public void processOneIsSuccessful() { @@ -138,7 +210,7 @@ public class CqlTemplateUnitTests { } /** - * @see DATACASS-286 + * @see DATACASS-286 */ @Test public void processOneThrowsIncorrectResultSetSizeDataAccessExceptionWhenNoRowsFound() { @@ -149,7 +221,6 @@ public class CqlTemplateUnitTests { when(mockResultSet.one()).thenReturn(null); try { - exception.expect(IncorrectResultSizeDataAccessException.class); exception.expectCause(is(nullValue(Throwable.class))); exception.expectMessage(containsString("expected 1, actual 0")); @@ -164,7 +235,7 @@ public class CqlTemplateUnitTests { } /** - * @see DATACASS-286 + * @see DATACASS-286 */ @Test public void processOneThrowsIncorrectResultSetSizeDataAccessExceptionWhenTooManyRowsFound() { @@ -177,7 +248,6 @@ public class CqlTemplateUnitTests { when(mockResultSet.isExhausted()).thenReturn(false); try { - exception.expect(IncorrectResultSizeDataAccessException.class); exception.expectCause(is(nullValue(Throwable.class))); exception.expectMessage("ResultSet size exceeds 1"); @@ -193,7 +263,7 @@ public class CqlTemplateUnitTests { } /** - * @see DATACASS-286 + * @see DATACASS-286 */ @Test public void processOnePassingNullResultSetThrowsIllegalArgumentException() { @@ -201,7 +271,6 @@ public class CqlTemplateUnitTests { RowMapper mockRowMapper = mock(RowMapper.class); try { - exception.expect(IllegalArgumentException.class); exception.expectCause(is(nullValue(Throwable.class))); @@ -212,7 +281,7 @@ public class CqlTemplateUnitTests { } /** - * @see DATACASS-286 + * @see DATACASS-286 */ @Test public void processOneWithRequiredTypeIsSuccessful() { @@ -242,7 +311,7 @@ public class CqlTemplateUnitTests { } /** - * @see DATACASS-286 + * @see DATACASS-286 */ @Test public void processOneWithRequiredTypeThrowsIncorrectResultSetSizeDataAccessExceptionWhenNoRowsFound() { @@ -252,7 +321,6 @@ public class CqlTemplateUnitTests { when(mockResultSet.one()).thenReturn(null); try { - exception.expect(IncorrectResultSizeDataAccessException.class); exception.expectCause(is(nullValue(Throwable.class))); exception.expectMessage(containsString("expected 1, actual 0")); @@ -266,7 +334,7 @@ public class CqlTemplateUnitTests { } /** - * @see DATACASS-286 + * @see DATACASS-286 */ @Test public void processOneWithRequiredTypeThrowsIncorrectResultSetSizeDataAccessExceptionWhenTooManyRowsFound() { @@ -278,7 +346,6 @@ public class CqlTemplateUnitTests { when(mockResultSet.isExhausted()).thenReturn(false); try { - exception.expect(IncorrectResultSizeDataAccessException.class); exception.expectCause(is(nullValue(Throwable.class))); exception.expectMessage(containsString("ResultSet size exceeds 1")); @@ -293,11 +360,10 @@ public class CqlTemplateUnitTests { } /** - * @see DATACASS-286 + * @see DATACASS-286 */ @Test public void processOneWithRequiredTypePassingNullResultSetThrowsIllegalArgumentException() { - exception.expect(IllegalArgumentException.class); exception.expectCause(is(nullValue(Throwable.class))); diff --git a/spring-cql/src/test/java/org/springframework/cassandra/core/support/EmptyResultSetUnitTests.java b/spring-cql/src/test/java/org/springframework/cassandra/core/support/EmptyResultSetUnitTests.java new file mode 100644 index 000000000..48007cbe5 --- /dev/null +++ b/spring-cql/src/test/java/org/springframework/cassandra/core/support/EmptyResultSetUnitTests.java @@ -0,0 +1,109 @@ +/* + * Copyright 2013-2016 the original author or authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cassandra.core.support; + +import static org.hamcrest.Matchers.*; +import static org.junit.Assert.*; +import static org.mockito.Mockito.*; + +import java.util.Collections; +import java.util.Iterator; + +import org.junit.Test; + +import com.datastax.driver.core.ExecutionInfo; +import com.datastax.driver.core.ResultSet; +import com.datastax.driver.core.Row; + +/** + * Test suite of test cases testing the contract and functionality of the {@link EmptyResultSet}. + * + * @author John Blum + * @see org.springframework.cassandra.core.support.EmptyResultSet + * @since 1.5.0 + */ +public class EmptyResultSetUnitTests { + + @Test + public void nullSafeResultSetReturnsGivenResultSet() { + ResultSet mockResultSet = mock(ResultSet.class); + ResultSet theResultSet = EmptyResultSet.nullSafeResultSet(mockResultSet); + + assertThat(theResultSet, is(sameInstance(mockResultSet))); + } + + @Test + public void nullSAfeResultSetReturnsEmptyResultSetForNull() { + ResultSet resultSet = EmptyResultSet.nullSafeResultSet(null); + + assertThat(resultSet, is(instanceOf(EmptyResultSet.class))); + } + + @Test + public void isExhaustedForEmptyResultIsTrue() { + assertThat(EmptyResultSet.nullSafeResultSet(null).isExhausted(), is(true)); + } + + @Test + public void isFullyFetchedForEmptyResultSetIsTrue() { + assertThat(EmptyResultSet.nullSafeResultSet(null).isFullyFetched(), is(true)); + } + + @Test + public void getAllExecutionInfoForEmptyResultSetIsEmptyList() { + assertThat(EmptyResultSet.nullSafeResultSet(null).getAllExecutionInfo(), + is(equalTo(Collections.emptyList()))); + } + + @Test(expected = UnsupportedOperationException.class) + public void getColumnDefinitionsForEmptyResultSetThrowsUnsupportedOperationException() { + EmptyResultSet.nullSafeResultSet(null).getColumnDefinitions(); + } + + @Test + public void getExecutionInfoForEmptyResultSetIsNull() { + assertThat(EmptyResultSet.nullSafeResultSet(null).getExecutionInfo(), is(nullValue(ExecutionInfo.class))); + } + + @Test + public void allForEmptyResultSetIsEmptyList() { + assertThat(EmptyResultSet.nullSafeResultSet(null).all(), is(equalTo(Collections.emptyList()))); + } + + @Test(expected = UnsupportedOperationException.class) + public void fetchMoreResultsFromEmptyResultSetThrowsUnsupportedOperationException() { + EmptyResultSet.nullSafeResultSet(null).fetchMoreResults(); + } + + @Test + public void iteratorForEmptyResultSetIsEmptyIterator() { + Iterator iterator = EmptyResultSet.nullSafeResultSet(null).iterator(); + + assertThat(iterator, is(notNullValue(Iterator.class))); + assertThat(iterator.hasNext(), is(false)); + } + + @Test + public void oneForEmptyResultSetIsNull() { + assertThat(EmptyResultSet.nullSafeResultSet(null).one(), is(nullValue(Row.class))); + } + + @Test(expected = UnsupportedOperationException.class) + public void wasAppliedOnEmptyResultSetThrowsUnsupportedOperationException() { + EmptyResultSet.nullSafeResultSet(null).wasApplied(); + } +} 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 74801f30c..07669c4b7 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 @@ -26,11 +26,11 @@ import org.springframework.cassandra.core.Cancellable; import org.springframework.cassandra.core.CqlTemplate; import org.springframework.cassandra.core.QueryForObjectListener; import org.springframework.cassandra.core.QueryOptions; -import org.springframework.cassandra.core.SessionCallback; +import org.springframework.cassandra.core.RowCallback; import org.springframework.cassandra.core.WriteOptions; import org.springframework.cassandra.core.cql.CqlIdentifier; +import org.springframework.cassandra.core.support.EmptyResultSet; import org.springframework.cassandra.core.util.CollectionUtils; -import org.springframework.dao.DataAccessException; import org.springframework.dao.DuplicateKeyException; import org.springframework.dao.InvalidDataAccessApiUsageException; import org.springframework.dao.support.PersistenceExceptionTranslator; @@ -77,24 +77,30 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation protected CassandraMappingContext mappingContext; /** - * Default Constructor for wiring in the required components later + * Default constructor used to wire in the required components later. */ public CassandraTemplate() {} /** - * Creates a new {@link} for the given {@link Session}. + * Creates a new {@link CassandraTemplate} for the given {@link Session}. * - * @param session must not be {@literal null}. + * @param session Cassandra {@link Session} connected to the Cassandra cluster instance; + * must not be {@literal null}. + * @see com.datastax.driver.core.Session */ public CassandraTemplate(Session session) { this(session, null); } /** - * Constructor if only session and converter are known at time of Template Creation + * Creates an instance of {@link CassandraTemplate} initialized with the given {@link Session} + * and {@link CassandraConverter}. * - * @param session must not be {@literal null}. - * @param converter must not be {@literal null}. + * @param session {@link Session} used to interact with Cassandra; must not be {@literal null}. + * @param converter {@link CassandraConverter} used to convert between Java and Cassandra types; + * must not be {@literal null}. + * @see org.springframework.data.cassandra.convert.CassandraConverter + * @see com.datastax.driver.core.Session */ public CassandraTemplate(Session session, CassandraConverter converter) { setSession(session); @@ -136,14 +142,22 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation } /** - * Returns the {@link CassandraMappingContext} - * - * @return the {@link CassandraMappingContext} + * @deprecated see {@link #getMappingContext()}. */ + @Deprecated public CassandraMappingContext getCassandraMappingContext() { return mappingContext; } + /** + * Returns the {@link CassandraMappingContext}. + * + * @return the {@link CassandraMappingContext}. + */ + public CassandraMappingContext getMappingContext() { + return mappingContext; + } + /* (non-Javadoc) * @see org.springframework.cassandra.support.CassandraAccessor#afterPropertiesSet() */ @@ -409,10 +423,6 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation return selectOne(select, entityClass); } - protected interface ClauseCallback { - void doWithClause(Clause clause); - } - @Deprecated protected void appendIdCriteria(ClauseCallback clauseCallback, CassandraPersistentEntity entity, Map id) { @@ -580,44 +590,34 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation return doUpdateAsync(entity, listener, options); } - protected List select(final String query, CassandraConverterRowCallback readRowCallback) { - - ResultSet resultSet = doExecute(new SessionCallback() { - - @Override - public ResultSet doInSession(Session session) throws DataAccessException { - return session.execute(query); - } - }); - - if (resultSet != null) { - List result = new ArrayList(); - - for (Row row : resultSet) { - result.add(readRowCallback.doWith(row)); - } - - return result; - } - - return null; + protected List select(String query, CassandraConverterRowCallback rowCallback) { + return processResultSet(doExecuteQueryReturnResultSet(query), rowCallback); } - /* (non-Javadoc) + protected List select(Select query, CassandraConverterRowCallback rowCallback) { + return processResultSet(doExecuteQueryReturnResultSet(query), rowCallback); + } + + private List processResultSet(ResultSet resultSet, RowCallback rowCallback) { + List result = new ArrayList(); + + for (Row row : EmptyResultSet.nullSafeResultSet(resultSet)) { + result.add(rowCallback.doWith(row)); + } + + return result; + } + + /* + * (non-Javadoc) * @see org.springframework.data.cassandra.core.CassandraOperations#stream(java.lang.String, java.lang.Class) */ - public Iterator stream(final String query, Class entityClass) { + public Iterator stream(String query, Class entityClass) { Assert.hasText(query, "Query must not be empty"); Assert.notNull(entityClass, "EntityClass must not be null"); - ResultSet resultSet = doExecute(new SessionCallback() { - - @Override - public ResultSet doInSession(Session session) throws DataAccessException { - return session.execute(logCql(query)); - } - }); + ResultSet resultSet = doExecuteQueryReturnResultSet(query); return (resultSet != null ? toIterator(resultSet, entityClass) : Collections.emptyIterator()); } @@ -633,29 +633,6 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation new CassandraConverterRowCallback(cassandraConverter, entityClass)); } - protected List select(final Select query, CassandraConverterRowCallback readRowCallback) { - - ResultSet resultSet = doExecute(new SessionCallback() { - - @Override - public ResultSet doInSession(Session session) throws DataAccessException { - return session.execute(query); - } - }); - - if (resultSet != null) { - List result = new ArrayList(); - - for (Row row : resultSet) { - result.add(readRowCallback.doWith(row)); - } - - return result; - } - - return null; - } - protected T selectOne(String query, CassandraConverterRowCallback rowCallback) { Iterator iterator = query(logCql(query)).iterator(); @@ -699,7 +676,7 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation // TODO: handle possible IndexOutOfBoundsException if the List of entities is empty protected void doBatchDelete(List entities, QueryOptions options) { execute(createDeleteBatchQuery(getTableName(entities.get(0).getClass()).toCql(), entities, options, - cassandraConverter)); + cassandraConverter)); } // TODO: handle possible IndexOutOfBoundsException if the List of entities is empty @@ -765,9 +742,9 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation protected List doBatchWrite(List entities, WriteOptions options, boolean insert) { - if (entities == null || entities.isEmpty()) { + if (CollectionUtils.isEmpty(entities)) { if (logger.isWarnEnabled()) { - logger.warn("no-op due to given null or empty list"); + logger.warn("no-op due to given null or empty List"); } return entities; @@ -826,7 +803,7 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation protected Cancellable doBatchWriteAsync(final List entities, final WriteListener listener, WriteOptions options, boolean insert) { - if (entities == null || entities.size() == 0) { + if (CollectionUtils.isEmpty(entities)) { if (logger.isWarnEnabled()) { logger.warn("no-op due to given null or empty list"); } @@ -866,9 +843,7 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation protected void doDelete(T entity, QueryOptions options) { Assert.notNull(entity, "Entity must not be null"); - Delete delete = createDeleteQuery(getTableName(entity.getClass()).toCql(), entity, options, cassandraConverter); - - execute(delete); + execute(createDeleteQuery(getTableName(entity.getClass()).toCql(), entity, options, cassandraConverter)); } protected Cancellable doDeleteAsync(final T entity, final DeletionListener listener, QueryOptions options) { @@ -895,9 +870,7 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation protected T doUpdate(T entity, WriteOptions options) { Assert.notNull(entity, "Entity must not be null"); - Update update = createUpdateQuery(getTableName(entity.getClass()).toCql(), entity, options, cassandraConverter); - - execute(update); + execute(createUpdateQuery(getTableName(entity.getClass()).toCql(), entity, options, cassandraConverter)); return entity; } @@ -906,15 +879,13 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation Assert.notNull(entity, "Entity must not be null"); - Update update = createUpdateQuery(getTableName(entity.getClass()).toCql(), entity, options, cassandraConverter); - AsynchronousQueryListener queryListener = (listener == null ? null : new AsynchronousQueryListener() { @Override @SuppressWarnings("unchecked") - public void onQueryComplete(ResultSetFuture rsf) { + public void onQueryComplete(ResultSetFuture resultSetFuture) { try { - rsf.getUninterruptibly(); + resultSetFuture.getUninterruptibly(); listener.onWriteComplete(Collections.singletonList(entity)); } catch (Exception x) { listener.onException(translateExceptionIfPossible(x)); @@ -922,7 +893,8 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation } }); - return executeAsynchronously(update, queryListener); + return executeAsynchronously(createUpdateQuery(getTableName(entity.getClass()).toCql(), entity, options, + cassandraConverter), queryListener); } /** @@ -941,10 +913,9 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation Assert.notNull(objectToUpdate, "Object to insert must not be null"); Assert.notNull(entityWriter, "EntityWriter must not be null"); - Insert insert = QueryBuilder.insertInto(tableName); + Insert insert = addWriteOptions(QueryBuilder.insertInto(tableName), options); entityWriter.write(objectToUpdate, insert); - CqlTemplate.addWriteOptions(insert, options); return insert; } @@ -992,10 +963,9 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation Assert.notNull(objectToUpdate, "Object to update must not be null"); Assert.notNull(entityWriter, "EntityWriter must not be null"); - Update update = QueryBuilder.update(tableName); + Update update = addWriteOptions(QueryBuilder.update(tableName), options); entityWriter.write(objectToUpdate, update); - CqlTemplate.addWriteOptions(update, options); return update; } @@ -1072,7 +1042,6 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation entityWriter.write(objectToDelete, where); - return delete; } @@ -1165,6 +1134,7 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation T result = new CassandraConverterRowCallback(cassandraConverter, entityClass).doWith(row); if (iterator.hasNext()) { + // TODO: throw IncorrectResultSetSizeDataAccessException instead throw new DuplicateKeyException(String.format( "found two or more results in query [%s]", query)); } @@ -1191,6 +1161,10 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation "Expected type String or Select; got type [%1$s] with value [%2$s]", query.getClass(), query)); } + protected interface ClauseCallback { + void doWithClause(Clause clause); + } + private static class ResultSetIteratorAdapter implements Iterator{ private final CassandraConverterRowCallback rowCallback; 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 new file mode 100644 index 000000000..d2be01a58 --- /dev/null +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/core/CassandraTemplateUnitTests.java @@ -0,0 +1,170 @@ +/* + * Copyright 2013-2016 the original author or authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.data.cassandra.core; + +import static org.hamcrest.Matchers.*; +import static org.junit.Assert.*; +import static org.mockito.Mockito.*; + +import java.util.Arrays; +import java.util.Collections; +import java.util.Iterator; +import java.util.List; + +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.Mock; +import org.mockito.runners.MockitoJUnitRunner; +import org.springframework.data.cassandra.convert.CassandraConverter; + +import com.datastax.driver.core.ResultSet; +import com.datastax.driver.core.Row; +import com.datastax.driver.core.Session; +import com.datastax.driver.core.querybuilder.Select; + +/** + * Test suite of test cases testing the contract and functionality of the {@link CassandraTemplate} class. + * + * @author John Blum + * @see org.springframework.data.cassandra.core.CassandraTemplate + * @since 1.5.0 + */ +@RunWith(MockitoJUnitRunner.class) +public class CassandraTemplateUnitTests { + + private CassandraTemplate template; + + @Mock + private Session mockSession; + + @Before + public void setup() { + template = new CassandraTemplate(mockSession); + } + + protected Iterator iterator(T... elements) { + return Collections.unmodifiableList(Arrays.asList(elements)).iterator(); + } + + protected Row mockRow(String name) { + return mock(Row.class, name); + } + protected CassandraConverterRowCallback newRollCallback(CassandraConverter converter, Class type) { + return new CassandraConverterRowCallback(converter, type); + } + + /** + * @see DATACASS-310 + */ + @Test + public void processResultSetHandlesResultSetRows() { + ResultSet mockResultSet = mock(ResultSet.class); + + Row mockRowOne = mockRow("MockRowOne"); + Row mockRowTwo = mockRow("MockRowTwo"); + Row mockRowThree = mockRow("MockRowThree"); + + CassandraConverter mockCassandraConverter = mock(CassandraConverter.class); + + when(mockSession.execute(eq("SELECT * FROM Test"))).thenReturn(mockResultSet); + when(mockResultSet.iterator()).thenReturn(iterator(mockRowOne, mockRowTwo, mockRowThree)); + when(mockCassandraConverter.read(eq(Integer.class), eq(mockRowOne))).thenReturn(1); + when(mockCassandraConverter.read(eq(Integer.class), eq(mockRowTwo))).thenReturn(2); + when(mockCassandraConverter.read(eq(Integer.class), eq(mockRowThree))).thenReturn(3); + + List results = template.select("SELECT * FROM Test", + newRollCallback(mockCassandraConverter, Integer.class)); + + assertThat(results, is(notNullValue(List.class))); + assertThat(results.size(), is(equalTo(3))); + assertThat(results.containsAll(Arrays.asList(1, 2, 3)), is(true)); + + verify(mockSession, times(1)).execute(eq("SELECT * FROM Test")); + verify(mockResultSet, times(1)).iterator(); + verify(mockCassandraConverter, times(1)).read(eq(Integer.class), eq(mockRowOne)); + verify(mockCassandraConverter, times(1)).read(eq(Integer.class), eq(mockRowTwo)); + verify(mockCassandraConverter, times(1)).read(eq(Integer.class), eq(mockRowThree)); + } + + /** + * @see DATACASS-310 + */ + @Test + public void processResultSetHandlesSingleElementResultSet() { + Select mockSelect = mock(Select.class); + ResultSet mockResultSet = mock(ResultSet.class); + Row mockRow = mock(Row.class); + CassandraConverter mockCassandraConverter = mock(CassandraConverter.class); + + when(mockSession.execute(eq(mockSelect))).thenReturn(mockResultSet); + when(mockResultSet.iterator()).thenReturn(iterator(mockRow)); + when(mockCassandraConverter.read(eq(String.class), eq(mockRow))).thenReturn("test"); + + List results = template.select(mockSelect, + newRollCallback(mockCassandraConverter, String.class)); + + assertThat(results, is(notNullValue(List.class))); + assertThat(results.size(), is(equalTo(1))); + assertThat(results, hasItem("test")); + + verify(mockSession, times(1)).execute(eq(mockSelect)); + verify(mockResultSet, times(1)).iterator(); + verify(mockCassandraConverter, times(1)).read(eq(String.class), eq(mockRow)); + } + + /** + * @see DATACASS-310 + */ + @Test + public void processResultSetHandlesEmptyResultSet() { + CassandraConverter mockCassandraConverter = mock(CassandraConverter.class); + ResultSet mockResultSet = mock(ResultSet.class); + + when(mockSession.execute(eq("SELECT * FROM Test"))).thenReturn(mockResultSet); + when(mockResultSet.iterator()).thenReturn(this.iterator()); + + List results = template.select("SELECT * FROM Test", + newRollCallback(mockCassandraConverter, Object.class)); + + assertThat(results, is(notNullValue(List.class))); + assertThat(results.isEmpty(), is(true)); + + verify(mockSession, times(1)).execute(eq("SELECT * FROM Test")); + verify(mockResultSet, times(1)).iterator(); + verifyZeroInteractions(mockCassandraConverter); + } + + /** + * @see DATACASS-310 + */ + @Test + public void processResultSetHandlesNullResultSet() { + CassandraConverter mockCassandraConverter = mock(CassandraConverter.class); + + when(mockSession.execute(anyString())).thenReturn(null); + + List results = template.select("SELECT * FROM Test", + newRollCallback(mockCassandraConverter, Object.class)); + + assertThat(results, is(notNullValue(List.class))); + assertThat(results.isEmpty(), is(true)); + + verify(mockSession, times(1)).execute(eq("SELECT * FROM Test")); + verifyZeroInteractions(mockCassandraConverter); + } +}