From 4578c9aab1614fd88700f8e356e3cf17b02a13c0 Mon Sep 17 00:00:00 2001 From: Bianca Stanciu Date: Wed, 22 Jul 2026 17:58:30 +0300 Subject: [PATCH 1/4] [CASSANALYTICS-183] - Mixed-case keyspace restore fails with Keyspace does not exist when quoteIdentifiers is set --- .../spark/data/QualifiedTableName.java | 16 +++ .../spark/data/QualifiedTableNameTest.java | 84 ++++++++++++ .../bulkwriter/CassandraClusterInfo.java | 7 +- .../CloudStorageDataTransferApiImpl.java | 20 +-- ...oordinatedCloudStorageDataTransferApi.java | 8 +- .../spark/data/CassandraDataLayer.java | 2 +- .../CloudStorageDataTransferApiImplTest.java | 120 ++++++++++++++++++ 7 files changed, 239 insertions(+), 18 deletions(-) create mode 100644 cassandra-analytics-common/src/test/java/org/apache/cassandra/spark/data/QualifiedTableNameTest.java create mode 100644 cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageDataTransferApiImplTest.java diff --git a/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/data/QualifiedTableName.java b/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/data/QualifiedTableName.java index da228d78b..36558c037 100644 --- a/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/data/QualifiedTableName.java +++ b/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/data/QualifiedTableName.java @@ -65,6 +65,14 @@ public String keyspace() return keyspace; } + /** + * @return the keyspace name, quoted with double quotes when {@code quoteIdentifiers} is set + */ + public String maybeQuotedKeyspace() + { + return maybeQuote(keyspace); + } + /** * @return the table name in Cassandra */ @@ -73,6 +81,14 @@ public String table() return table; } + /** + * @return the table name, quoted with double quotes when {@code quoteIdentifiers} is set + */ + public String maybeQuotedTable() + { + return maybeQuote(table); + } + /** * @return the identifiers should be quoted */ diff --git a/cassandra-analytics-common/src/test/java/org/apache/cassandra/spark/data/QualifiedTableNameTest.java b/cassandra-analytics-common/src/test/java/org/apache/cassandra/spark/data/QualifiedTableNameTest.java new file mode 100644 index 000000000..ac27beb79 --- /dev/null +++ b/cassandra-analytics-common/src/test/java/org/apache/cassandra/spark/data/QualifiedTableNameTest.java @@ -0,0 +1,84 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you 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.apache.cassandra.spark.data; + +import org.junit.jupiter.api.Test; + +import static org.assertj.core.api.Assertions.assertThat; + +class QualifiedTableNameTest +{ + @Test + void testUnquotedIdentifiersReturnedAsIs() + { + QualifiedTableName name = new QualifiedTableName("mykeyspace", "mytable", false); + assertThat(name.keyspace()).isEqualTo("mykeyspace"); + assertThat(name.table()).isEqualTo("mytable"); + assertThat(name.maybeQuotedKeyspace()).isEqualTo("mykeyspace"); + assertThat(name.maybeQuotedTable()).isEqualTo("mytable"); + } + + @Test + void testQuotedIdentifiersWrappedInDoubleQuotes() + { + QualifiedTableName name = new QualifiedTableName("AAMBackend", "UserBackend", true); + assertThat(name.maybeQuotedKeyspace()).isEqualTo("\"AAMBackend\""); + assertThat(name.maybeQuotedTable()).isEqualTo("\"UserBackend\""); + } + + @Test + void testRawAccessorsUnaffectedByQuoteFlag() + { + QualifiedTableName name = new QualifiedTableName("AAMBackend", "UserBackend", true); + assertThat(name.keyspace()).isEqualTo("AAMBackend"); + assertThat(name.table()).isEqualTo("UserBackend"); + } + + @Test + void testToStringQuotesBothWhenFlagSet() + { + QualifiedTableName name = new QualifiedTableName("AAMBackend", "UserBackend", true); + assertThat(name.toString()).isEqualTo("\"AAMBackend\".\"UserBackend\""); + } + + @Test + void testToStringUnquotedWhenFlagNotSet() + { + QualifiedTableName name = new QualifiedTableName("mykeyspace", "mytable", false); + assertThat(name.toString()).isEqualTo("mykeyspace.mytable"); + } + + @Test + void testDefaultConstructorDoesNotQuote() + { + QualifiedTableName name = new QualifiedTableName("AAMBackend", "UserBackend"); + assertThat(name.quoteIdentifiers()).isFalse(); + assertThat(name.maybeQuotedKeyspace()).isEqualTo("AAMBackend"); + assertThat(name.maybeQuotedTable()).isEqualTo("UserBackend"); + } + + @Test + void testReservedWordIdentifiersQuoted() + { + QualifiedTableName name = new QualifiedTableName("keyspace", "table", true); + assertThat(name.maybeQuotedKeyspace()).isEqualTo("\"keyspace\""); + assertThat(name.maybeQuotedTable()).isEqualTo("\"table\""); + } +} diff --git a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/CassandraClusterInfo.java b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/CassandraClusterInfo.java index 1b1a40179..5d157974a 100644 --- a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/CassandraClusterInfo.java +++ b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/CassandraClusterInfo.java @@ -289,11 +289,12 @@ protected String getCurrentKeyspaceSchema() throws Exception private TokenRangeReplicasResponse getTokenRangesAndReplicaSets() { CassandraContext context = getCassandraContext(); + String quotedKeyspace = maybeQuotedIdentifier(bridge(), conf.quoteIdentifiers, conf.keyspace); try { long start = System.nanoTime(); TokenRangeReplicasResponse response = context.getSidecarClient() - .tokenRangeReplicas(new ArrayList<>(context.getCluster()), conf.keyspace) + .tokenRangeReplicas(new ArrayList<>(context.getCluster()), quotedKeyspace) .get(); long elapsedTimeNanos = System.nanoTime() - start; LOGGER.info("Retrieved token ranges for {} instances in {} milliseconds", @@ -303,8 +304,8 @@ private TokenRangeReplicasResponse getTokenRangesAndReplicaSets() } catch (ExecutionException | InterruptedException exception) { - LOGGER.error("Failed to get token ranges for keyspace {}", conf.keyspace, exception); - throw new SidecarApiCallException("Failed to get token ranges for keyspace" + conf.keyspace, exception); + LOGGER.error("Failed to get token ranges for keyspace {}", quotedKeyspace, exception); + throw new SidecarApiCallException("Failed to get token ranges for keyspace " + quotedKeyspace, exception); } } diff --git a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageDataTransferApiImpl.java b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageDataTransferApiImpl.java index be218365f..f72a7a271 100644 --- a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageDataTransferApiImpl.java +++ b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageDataTransferApiImpl.java @@ -103,8 +103,8 @@ public void createRestoreJob(CreateRestoreJobRequestPayload createRestoreJobRequ try { QualifiedTableName qualifiedTableName = jobInfo.qualifiedTableName(); - sidecarClient.createRestoreJob(qualifiedTableName.keyspace(), - qualifiedTableName.table(), + sidecarClient.createRestoreJob(qualifiedTableName.maybeQuotedKeyspace(), + qualifiedTableName.maybeQuotedTable(), createRestoreJobRequestPayload).get(); } catch (Exception exception) @@ -121,8 +121,8 @@ public RestoreJobSummaryResponsePayload restoreJobSummary() throws SidecarApiCal try { QualifiedTableName qualifiedTableName = jobInfo.qualifiedTableName(); - return sidecarClient.restoreJobSummary(qualifiedTableName.keyspace(), - qualifiedTableName.table(), + return sidecarClient.restoreJobSummary(qualifiedTableName.maybeQuotedKeyspace(), + qualifiedTableName.maybeQuotedTable(), jobId).get(); } catch (Exception exception) @@ -168,8 +168,8 @@ public void updateRestoreJob(UpdateRestoreJobRequestPayload updateRestoreJobRequ { LOGGER.info("Updating the restore job. clusterId={} restoreJobId={}", clusterId, jobId); QualifiedTableName qualifiedTableName = jobInfo.qualifiedTableName(); - sidecarClient.updateRestoreJob(qualifiedTableName.keyspace(), - qualifiedTableName.table(), + sidecarClient.updateRestoreJob(qualifiedTableName.maybeQuotedKeyspace(), + qualifiedTableName.maybeQuotedTable(), jobId, updateRestoreJobRequestPayload).get(); } @@ -188,8 +188,8 @@ public void abortRestoreJob() throws SidecarApiCallException { LOGGER.info("Abort job. clusterId={} restoreJobId={}", clusterId, jobId); QualifiedTableName qualifiedTableName = jobInfo.qualifiedTableName(); - sidecarClient.abortRestoreJob(qualifiedTableName.keyspace(), - qualifiedTableName.table(), + sidecarClient.abortRestoreJob(qualifiedTableName.maybeQuotedKeyspace(), + qualifiedTableName.maybeQuotedTable(), jobId).get(); } catch (Exception exception) @@ -207,8 +207,8 @@ private CompletableFuture createRestoreSliceWithCustomRetry(SidecarInstanc RetryPolicy retryPolicy) { QualifiedTableName qualifiedTableName = jobInfo.qualifiedTableName(); - CreateRestoreJobSliceRequest request = new CreateRestoreJobSliceRequest(qualifiedTableName.keyspace(), - qualifiedTableName.table(), + CreateRestoreJobSliceRequest request = new CreateRestoreJobSliceRequest(qualifiedTableName.maybeQuotedKeyspace(), + qualifiedTableName.maybeQuotedTable(), jobInfo.getRestoreJobId(clusterId), createSliceRequestPayload); return sidecarClient.executeRequestAsync(sidecarClient.requestBuilder() diff --git a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/coordinated/CoordinatedCloudStorageDataTransferApi.java b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/coordinated/CoordinatedCloudStorageDataTransferApi.java index 4978a1eb1..0a1b6dffa 100644 --- a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/coordinated/CoordinatedCloudStorageDataTransferApi.java +++ b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/coordinated/CoordinatedCloudStorageDataTransferApi.java @@ -169,8 +169,8 @@ private void createRestoreSliceInternal(String clusterId, SidecarClient sidecarClient = dataTransferApi.sidecarClient(); QualifiedTableName qualifiedTableName = jobInfo.qualifiedTableName(); UUID restoreJobId = jobInfo.getRestoreJobId(clusterId); - CreateRestoreJobSliceRequest request = new CreateRestoreJobSliceRequest(qualifiedTableName.keyspace(), - qualifiedTableName.table(), + CreateRestoreJobSliceRequest request = new CreateRestoreJobSliceRequest(qualifiedTableName.maybeQuotedKeyspace(), + qualifiedTableName.maybeQuotedTable(), restoreJobId, createSliceRequestPayload); RetryPolicy retryPolicy = new CloudStorageDataTransferApiImpl.ExecutorCreateSliceRetryPolicy(sidecarClient); @@ -197,8 +197,8 @@ private void restoreJobProgressInternal(String clusterId, JobInfo jobInfo = dataTransferApi.jobInfo(); QualifiedTableName qualifiedTableName = jobInfo.qualifiedTableName(); UUID restoreJobId = jobInfo.getRestoreJobId(clusterId); - RestoreJobProgressRequestParams requestParams = new RestoreJobProgressRequestParams(qualifiedTableName.keyspace(), - qualifiedTableName.table(), + RestoreJobProgressRequestParams requestParams = new RestoreJobProgressRequestParams(qualifiedTableName.maybeQuotedKeyspace(), + qualifiedTableName.maybeQuotedTable(), restoreJobId, fetchPolicy); RestoreJobProgressResponsePayload jobProgress = restoreJobProgress(dataTransferApi, requestParams); diff --git a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/data/CassandraDataLayer.java b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/data/CassandraDataLayer.java index ff7e0b67d..d70ce8728 100644 --- a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/data/CassandraDataLayer.java +++ b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/data/CassandraDataLayer.java @@ -1189,7 +1189,7 @@ protected Sizing getSizing(CompletableFuture ringFuture, ReplicationFactor replicationFactor, ClientConfig options) { - return SizingFactory.create(replicationFactor, options, consistencyLevel, keyspace, table, datacenter, sidecar, sidecarPort, ringFuture); + return SizingFactory.create(replicationFactor, options, consistencyLevel, maybeQuotedKeyspace, maybeQuotedTable, datacenter, sidecar, sidecarPort, ringFuture); } protected void await(CountDownLatch latch) diff --git a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageDataTransferApiImplTest.java b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageDataTransferApiImplTest.java new file mode 100644 index 000000000..21ab733d5 --- /dev/null +++ b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageDataTransferApiImplTest.java @@ -0,0 +1,120 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you 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.apache.cassandra.spark.bulkwriter.cloudstorage; + +import java.util.UUID; +import java.util.concurrent.CompletableFuture; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import o.a.c.sidecar.client.shaded.common.request.data.CreateRestoreJobRequestPayload; +import o.a.c.sidecar.client.shaded.common.request.data.UpdateRestoreJobRequestPayload; +import o.a.c.sidecar.client.shaded.common.response.data.RestoreJobSummaryResponsePayload; +import o.a.c.sidecar.client.shaded.client.SidecarClient; +import org.apache.cassandra.spark.bulkwriter.JobInfo; +import org.apache.cassandra.spark.data.QualifiedTableName; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +class CloudStorageDataTransferApiImplTest +{ + private static final String QUOTED_KEYSPACE = "\"AAMBackend\""; + private static final String QUOTED_TABLE = "\"UserBackend\""; + private static final UUID JOB_ID = UUID.randomUUID(); + + private SidecarClient sidecarClient; + private JobInfo jobInfo; + private CloudStorageDataTransferApiImpl api; + + @BeforeEach + void setup() + { + sidecarClient = mock(SidecarClient.class); + jobInfo = mock(JobInfo.class); + when(jobInfo.qualifiedTableName()).thenReturn(new QualifiedTableName("AAMBackend", "UserBackend", true)); + when(jobInfo.getRestoreJobId()).thenReturn(JOB_ID); + when(jobInfo.getRestoreJobId(null)).thenReturn(JOB_ID); + api = new CloudStorageDataTransferApiImpl(jobInfo, sidecarClient, mock(StorageClient.class), null); + } + + @Test + void testCreateRestoreJobUsesQuotedIdentifiers() throws Exception + { + when(sidecarClient.createRestoreJob(eq(QUOTED_KEYSPACE), eq(QUOTED_TABLE), any())) + .thenReturn(CompletableFuture.completedFuture(null)); + + api.createRestoreJob(mock(CreateRestoreJobRequestPayload.class)); + + verify(sidecarClient).createRestoreJob(eq(QUOTED_KEYSPACE), eq(QUOTED_TABLE), any()); + } + + @Test + void testRestoreJobSummaryUsesQuotedIdentifiers() throws Exception + { + RestoreJobSummaryResponsePayload response = mock(RestoreJobSummaryResponsePayload.class); + when(sidecarClient.restoreJobSummary(eq(QUOTED_KEYSPACE), eq(QUOTED_TABLE), eq(JOB_ID))) + .thenReturn(CompletableFuture.completedFuture(response)); + + RestoreJobSummaryResponsePayload result = api.restoreJobSummary(); + + assertThat(result).isSameAs(response); + verify(sidecarClient).restoreJobSummary(eq(QUOTED_KEYSPACE), eq(QUOTED_TABLE), eq(JOB_ID)); + } + + @Test + void testUpdateRestoreJobUsesQuotedIdentifiers() throws Exception + { + when(sidecarClient.updateRestoreJob(eq(QUOTED_KEYSPACE), eq(QUOTED_TABLE), eq(JOB_ID), any())) + .thenReturn(CompletableFuture.completedFuture(null)); + + api.updateRestoreJob(mock(UpdateRestoreJobRequestPayload.class)); + + verify(sidecarClient).updateRestoreJob(eq(QUOTED_KEYSPACE), eq(QUOTED_TABLE), eq(JOB_ID), any()); + } + + @Test + void testAbortRestoreJobUsesQuotedIdentifiers() throws Exception + { + when(sidecarClient.abortRestoreJob(eq(QUOTED_KEYSPACE), eq(QUOTED_TABLE), eq(JOB_ID))) + .thenReturn(CompletableFuture.completedFuture(null)); + + api.abortRestoreJob(); + + verify(sidecarClient).abortRestoreJob(eq(QUOTED_KEYSPACE), eq(QUOTED_TABLE), eq(JOB_ID)); + } + + @Test + void testUnquotedIdentifiersPassedAsIsWhenFlagNotSet() throws Exception + { + when(jobInfo.qualifiedTableName()).thenReturn(new QualifiedTableName("aambackend", "userbackend", false)); + when(sidecarClient.createRestoreJob(eq("aambackend"), eq("userbackend"), any())) + .thenReturn(CompletableFuture.completedFuture(null)); + + api.createRestoreJob(mock(CreateRestoreJobRequestPayload.class)); + + verify(sidecarClient).createRestoreJob(eq("aambackend"), eq("userbackend"), any()); + } +} From ea64e1c2042af12ca62b7cd2b6b1a54f6c522513 Mon Sep 17 00:00:00 2001 From: Bianca Stanciu Date: Wed, 22 Jul 2026 18:02:58 +0300 Subject: [PATCH 2/4] Fixing the tests --- .../spark/data/QualifiedTableNameTest.java | 22 +++++++++---------- .../CloudStorageDataTransferApiImplTest.java | 12 +++++----- 2 files changed, 17 insertions(+), 17 deletions(-) diff --git a/cassandra-analytics-common/src/test/java/org/apache/cassandra/spark/data/QualifiedTableNameTest.java b/cassandra-analytics-common/src/test/java/org/apache/cassandra/spark/data/QualifiedTableNameTest.java index ac27beb79..36ed62a92 100644 --- a/cassandra-analytics-common/src/test/java/org/apache/cassandra/spark/data/QualifiedTableNameTest.java +++ b/cassandra-analytics-common/src/test/java/org/apache/cassandra/spark/data/QualifiedTableNameTest.java @@ -38,24 +38,24 @@ void testUnquotedIdentifiersReturnedAsIs() @Test void testQuotedIdentifiersWrappedInDoubleQuotes() { - QualifiedTableName name = new QualifiedTableName("AAMBackend", "UserBackend", true); - assertThat(name.maybeQuotedKeyspace()).isEqualTo("\"AAMBackend\""); - assertThat(name.maybeQuotedTable()).isEqualTo("\"UserBackend\""); + QualifiedTableName name = new QualifiedTableName("MyKeyspace", "MyTable", true); + assertThat(name.maybeQuotedKeyspace()).isEqualTo("\"MyKeyspace\""); + assertThat(name.maybeQuotedTable()).isEqualTo("\"MyTable\""); } @Test void testRawAccessorsUnaffectedByQuoteFlag() { - QualifiedTableName name = new QualifiedTableName("AAMBackend", "UserBackend", true); - assertThat(name.keyspace()).isEqualTo("AAMBackend"); - assertThat(name.table()).isEqualTo("UserBackend"); + QualifiedTableName name = new QualifiedTableName("MyKeyspace", "MyTable", true); + assertThat(name.keyspace()).isEqualTo("MyKeyspace"); + assertThat(name.table()).isEqualTo("MyTable"); } @Test void testToStringQuotesBothWhenFlagSet() { - QualifiedTableName name = new QualifiedTableName("AAMBackend", "UserBackend", true); - assertThat(name.toString()).isEqualTo("\"AAMBackend\".\"UserBackend\""); + QualifiedTableName name = new QualifiedTableName("MyKeyspace", "MyTable", true); + assertThat(name.toString()).isEqualTo("\"MyKeyspace\".\"MyTable\""); } @Test @@ -68,10 +68,10 @@ void testToStringUnquotedWhenFlagNotSet() @Test void testDefaultConstructorDoesNotQuote() { - QualifiedTableName name = new QualifiedTableName("AAMBackend", "UserBackend"); + QualifiedTableName name = new QualifiedTableName("MyKeyspace", "MyTable"); assertThat(name.quoteIdentifiers()).isFalse(); - assertThat(name.maybeQuotedKeyspace()).isEqualTo("AAMBackend"); - assertThat(name.maybeQuotedTable()).isEqualTo("UserBackend"); + assertThat(name.maybeQuotedKeyspace()).isEqualTo("MyKeyspace"); + assertThat(name.maybeQuotedTable()).isEqualTo("MyTable"); } @Test diff --git a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageDataTransferApiImplTest.java b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageDataTransferApiImplTest.java index 21ab733d5..a9492085c 100644 --- a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageDataTransferApiImplTest.java +++ b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageDataTransferApiImplTest.java @@ -41,8 +41,8 @@ class CloudStorageDataTransferApiImplTest { - private static final String QUOTED_KEYSPACE = "\"AAMBackend\""; - private static final String QUOTED_TABLE = "\"UserBackend\""; + private static final String QUOTED_KEYSPACE = "\"MyKeyspace\""; + private static final String QUOTED_TABLE = "\"MyTable\""; private static final UUID JOB_ID = UUID.randomUUID(); private SidecarClient sidecarClient; @@ -54,7 +54,7 @@ void setup() { sidecarClient = mock(SidecarClient.class); jobInfo = mock(JobInfo.class); - when(jobInfo.qualifiedTableName()).thenReturn(new QualifiedTableName("AAMBackend", "UserBackend", true)); + when(jobInfo.qualifiedTableName()).thenReturn(new QualifiedTableName("MyKeyspace", "MyTable", true)); when(jobInfo.getRestoreJobId()).thenReturn(JOB_ID); when(jobInfo.getRestoreJobId(null)).thenReturn(JOB_ID); api = new CloudStorageDataTransferApiImpl(jobInfo, sidecarClient, mock(StorageClient.class), null); @@ -109,12 +109,12 @@ void testAbortRestoreJobUsesQuotedIdentifiers() throws Exception @Test void testUnquotedIdentifiersPassedAsIsWhenFlagNotSet() throws Exception { - when(jobInfo.qualifiedTableName()).thenReturn(new QualifiedTableName("aambackend", "userbackend", false)); - when(sidecarClient.createRestoreJob(eq("aambackend"), eq("userbackend"), any())) + when(jobInfo.qualifiedTableName()).thenReturn(new QualifiedTableName("mykeyspace", "mytable", false)); + when(sidecarClient.createRestoreJob(eq("mykeyspace"), eq("mytable"), any())) .thenReturn(CompletableFuture.completedFuture(null)); api.createRestoreJob(mock(CreateRestoreJobRequestPayload.class)); - verify(sidecarClient).createRestoreJob(eq("aambackend"), eq("userbackend"), any()); + verify(sidecarClient).createRestoreJob(eq("mykeyspace"), eq("mytable"), any()); } } From e59fe5b4a9ee3c444ff046f38ffcb799d4445916 Mon Sep 17 00:00:00 2001 From: Bianca Stanciu Date: Thu, 23 Jul 2026 15:55:43 +0300 Subject: [PATCH 3/4] Fix checkstyle violation and add changelog entry --- CHANGES.txt | 1 + .../org/apache/cassandra/spark/data/CassandraDataLayer.java | 4 +++- 2 files changed, 4 insertions(+), 1 deletion(-) diff --git a/CHANGES.txt b/CHANGES.txt index 811417a7e..de9f9f770 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,5 +1,6 @@ 0.5.0 ----- + * Mixed-case keyspace restore fails with keyspace does not exist when quoteIdentifiers is set (CASSANALYTICS-183) * CdcState.ReplicaCountSerializer map-size overflow corrupts persisted CDC state (CASSANALYTICS-184) * SSTable-version-based bridge determination (CASSANALYTICS-24) * Upgrade sidecar version to 0.4.0 (CASSANALYTICS-176) diff --git a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/data/CassandraDataLayer.java b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/data/CassandraDataLayer.java index d70ce8728..fcffabd12 100644 --- a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/data/CassandraDataLayer.java +++ b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/data/CassandraDataLayer.java @@ -1189,7 +1189,9 @@ protected Sizing getSizing(CompletableFuture ringFuture, ReplicationFactor replicationFactor, ClientConfig options) { - return SizingFactory.create(replicationFactor, options, consistencyLevel, maybeQuotedKeyspace, maybeQuotedTable, datacenter, sidecar, sidecarPort, ringFuture); + return SizingFactory.create(replicationFactor, options, consistencyLevel, + maybeQuotedKeyspace, maybeQuotedTable, datacenter, + sidecar, sidecarPort, ringFuture); } protected void await(CountDownLatch latch) From cc23f9dc4204c5efa6eb56c623c03def7339b3c9 Mon Sep 17 00:00:00 2001 From: Bianca Stanciu Date: Fri, 2 Oct 2026 13:58:10 +0300 Subject: [PATCH 4/4] CASSANALYTICS-183: Cover quoted identifiers across Sidecar APIs Rename the conditional keyspace value and verify quoted names in restore, coordinated, topology, and sizing requests. --- .../bulkwriter/CassandraClusterInfo.java | 8 +- .../bulkwriter/CassandraClusterInfoTest.java | 54 +++++++++ .../CloudStorageDataTransferApiImplTest.java | 26 +++++ ...inatedCloudStorageDataTransferApiTest.java | 104 ++++++++++++++++++ .../spark/data/CassandraDataLayerTests.java | 42 +++++++ 5 files changed, 230 insertions(+), 4 deletions(-) create mode 100644 cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/coordinated/CoordinatedCloudStorageDataTransferApiTest.java diff --git a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/CassandraClusterInfo.java b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/CassandraClusterInfo.java index ca0028bdc..406828f34 100644 --- a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/CassandraClusterInfo.java +++ b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/CassandraClusterInfo.java @@ -289,12 +289,12 @@ protected String getCurrentKeyspaceSchema() throws Exception private TokenRangeReplicasResponse getTokenRangesAndReplicaSets() { CassandraContext context = getCassandraContext(); - String quotedKeyspace = maybeQuotedIdentifier(bridge(), conf.quoteIdentifiers, conf.keyspace); + String maybeQuotedKeyspace = maybeQuotedIdentifier(bridge(), conf.quoteIdentifiers, conf.keyspace); try { long start = System.nanoTime(); TokenRangeReplicasResponse response = context.getSidecarClient() - .tokenRangeReplicas(new ArrayList<>(context.getCluster()), quotedKeyspace) + .tokenRangeReplicas(new ArrayList<>(context.getCluster()), maybeQuotedKeyspace) .get(); long elapsedTimeNanos = System.nanoTime() - start; LOGGER.info("Retrieved token ranges for {} instances in {} milliseconds", @@ -304,8 +304,8 @@ private TokenRangeReplicasResponse getTokenRangesAndReplicaSets() } catch (ExecutionException | InterruptedException exception) { - LOGGER.error("Failed to get token ranges for keyspace {}", quotedKeyspace, exception); - throw new SidecarApiCallException("Failed to get token ranges for keyspace " + quotedKeyspace, exception); + LOGGER.error("Failed to get token ranges for keyspace {}", maybeQuotedKeyspace, exception); + throw new SidecarApiCallException("Failed to get token ranges for keyspace " + maybeQuotedKeyspace, exception); } } diff --git a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/CassandraClusterInfoTest.java b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/CassandraClusterInfoTest.java index 99b06712f..60027ceaf 100644 --- a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/CassandraClusterInfoTest.java +++ b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/CassandraClusterInfoTest.java @@ -24,12 +24,15 @@ import java.util.ArrayList; import java.util.Collections; import java.util.List; +import java.util.Map; +import java.util.TreeMap; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; import java.util.stream.Stream; import com.google.common.collect.ImmutableMap; import com.google.common.util.concurrent.Uninterruptibles; +import org.apache.spark.SparkConf; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; import org.junit.jupiter.params.ParameterizedTest; @@ -37,20 +40,71 @@ import org.junit.jupiter.params.provider.MethodSource; import o.a.c.sidecar.client.shaded.common.response.NodeSettings; +import o.a.c.sidecar.client.shaded.common.response.SchemaResponse; import o.a.c.sidecar.client.shaded.common.response.TimeSkewResponse; +import o.a.c.sidecar.client.shaded.common.response.TokenRangeReplicasResponse; +import o.a.c.sidecar.client.shaded.client.SidecarClient; +import org.apache.cassandra.bridge.CassandraBridge; import org.apache.cassandra.spark.bulkwriter.token.TokenRangeMapping; import org.apache.cassandra.spark.exception.TimeSkewTooLargeException; import static org.apache.cassandra.spark.TestUtils.range; +import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatNoException; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.RETURNS_DEEP_STUBS; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; public class CassandraClusterInfoTest { + @Test + void testQuotedKeyspacePassedToSchemaAndTokenRanges() throws Exception + { + Map options = new TreeMap<>(String.CASE_INSENSITIVE_ORDER); + options.put(WriterOptions.SIDECAR_INSTANCES.name(), "127.0.0.1"); + options.put(WriterOptions.KEYSPACE.name(), "MyKeyspace"); + options.put(WriterOptions.TABLE.name(), "MyTable"); + options.put(WriterOptions.QUOTE_IDENTIFIERS.name(), "true"); + BulkSparkConf conf = new BulkSparkConf(new SparkConf(), options); + CassandraBridge bridge = mock(CassandraBridge.class); + when(bridge.maybeQuoteIdentifier("MyKeyspace")).thenReturn("\"MyKeyspace\""); + CassandraContext context = mock(CassandraContext.class); + SidecarClient sidecar = mock(SidecarClient.class); + when(context.getSidecarClient()).thenReturn(sidecar); + when(context.getCluster()).thenReturn(Collections.emptySet()); + SchemaResponse schema = mock(SchemaResponse.class); + when(schema.schema()).thenReturn("schema"); + when(sidecar.schema("\"MyKeyspace\"")).thenReturn(CompletableFuture.completedFuture(schema)); + CompletableFuture failed = new CompletableFuture<>(); + failed.completeExceptionally(new IllegalStateException("test failure")); + when(sidecar.tokenRangeReplicas(any(), eq("\"MyKeyspace\""))).thenReturn(failed); + + try (CassandraClusterInfo ci = new CassandraClusterInfo(conf) + { + @Override + protected CassandraContext buildCassandraContext() + { + return context; + } + + @Override + protected CassandraBridge bridge() + { + return bridge; + } + }) + { + assertThat(ci.getCurrentKeyspaceSchema()).isEqualTo("schema"); + assertThatThrownBy(() -> ci.getTokenRangeMapping(false)).hasMessageContaining("Unable to initialize ring information"); + verify(sidecar).schema("\"MyKeyspace\""); + verify(sidecar).tokenRangeReplicas(any(), eq("\"MyKeyspace\"")); + } + } + @Test void testTimeSkewAcceptable() { diff --git a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageDataTransferApiImplTest.java b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageDataTransferApiImplTest.java index a9492085c..ec7ed17ff 100644 --- a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageDataTransferApiImplTest.java +++ b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageDataTransferApiImplTest.java @@ -24,11 +24,15 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; import o.a.c.sidecar.client.shaded.common.request.data.CreateRestoreJobRequestPayload; +import o.a.c.sidecar.client.shaded.common.request.data.CreateSliceRequestPayload; import o.a.c.sidecar.client.shaded.common.request.data.UpdateRestoreJobRequestPayload; import o.a.c.sidecar.client.shaded.common.response.data.RestoreJobSummaryResponsePayload; +import o.a.c.sidecar.client.shaded.client.RequestContext; import o.a.c.sidecar.client.shaded.client.SidecarClient; +import o.a.c.sidecar.client.shaded.client.SidecarInstance; import org.apache.cassandra.spark.bulkwriter.JobInfo; import org.apache.cassandra.spark.data.QualifiedTableName; @@ -36,6 +40,7 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -106,6 +111,27 @@ void testAbortRestoreJobUsesQuotedIdentifiers() throws Exception verify(sidecarClient).abortRestoreJob(eq(QUOTED_KEYSPACE), eq(QUOTED_TABLE), eq(JOB_ID)); } + @Test + void testRestoreSliceUsesQuotedIdentifiers() throws Exception + { + when(sidecarClient.requestBuilder()).thenReturn(new RequestContext.Builder()); + when(sidecarClient.executeRequestAsync(any(RequestContext.class))) + .thenReturn(CompletableFuture.completedFuture(null)); + SidecarInstance instance = mock(SidecarInstance.class); + CreateSliceRequestPayload payload = mock(CreateSliceRequestPayload.class); + + api.createRestoreSliceFromExecutor(instance, payload); + api.createRestoreSliceFromDriver(instance, payload).get(); + + ArgumentCaptor requests = ArgumentCaptor.forClass(RequestContext.class); + verify(sidecarClient, times(2)).executeRequestAsync(requests.capture()); + for (RequestContext request : requests.getAllValues()) + { + assertThat(request.request().requestURI()).contains(QUOTED_KEYSPACE, QUOTED_TABLE, JOB_ID.toString()); + assertThat(request.request().requestBody()).isSameAs(payload); + } + } + @Test void testUnquotedIdentifiersPassedAsIsWhenFlagNotSet() throws Exception { diff --git a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/coordinated/CoordinatedCloudStorageDataTransferApiTest.java b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/coordinated/CoordinatedCloudStorageDataTransferApiTest.java new file mode 100644 index 000000000..df114bcef --- /dev/null +++ b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/coordinated/CoordinatedCloudStorageDataTransferApiTest.java @@ -0,0 +1,104 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you 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.apache.cassandra.spark.bulkwriter.cloudstorage.coordinated; + +import java.util.Collections; +import java.util.UUID; +import java.util.concurrent.CompletableFuture; + +import com.google.common.util.concurrent.RateLimiter; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; + +import o.a.c.sidecar.client.shaded.common.data.RestoreJobProgressFetchPolicy; +import o.a.c.sidecar.client.shaded.common.request.data.CreateSliceRequestPayload; +import o.a.c.sidecar.client.shaded.common.request.data.RestoreJobProgressRequestParams; +import o.a.c.sidecar.client.shaded.common.response.data.RestoreJobProgressResponsePayload; +import o.a.c.sidecar.client.shaded.client.RequestContext; +import o.a.c.sidecar.client.shaded.client.SidecarClient; +import org.apache.cassandra.spark.bulkwriter.JobInfo; +import org.apache.cassandra.spark.bulkwriter.cloudstorage.CloudStorageDataTransferApiImpl; +import org.apache.cassandra.spark.data.QualifiedTableName; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +class CoordinatedCloudStorageDataTransferApiTest +{ + private static final String CLUSTER_ID = "cluster"; + private static final UUID JOB_ID = UUID.randomUUID(); + + private SidecarClient sidecarClient; + private CoordinatedCloudStorageDataTransferApi api; + + @BeforeEach + void setup() + { + sidecarClient = mock(SidecarClient.class); + JobInfo jobInfo = mock(JobInfo.class); + when(jobInfo.qualifiedTableName()).thenReturn(new QualifiedTableName("MyKeyspace", "MyTable", true)); + when(jobInfo.getRestoreJobId(CLUSTER_ID)).thenReturn(JOB_ID); + CloudStorageDataTransferApiImpl delegate = mock(CloudStorageDataTransferApiImpl.class); + when(delegate.jobInfo()).thenReturn(jobInfo); + when(delegate.sidecarClient()).thenReturn(sidecarClient); + api = new CoordinatedCloudStorageDataTransferApi(mock(RateLimiter.class), + Collections.singletonMap(CLUSTER_ID, delegate)); + } + + @Test + void testCreateRestoreSliceUsesQuotedIdentifiers() + { + when(sidecarClient.requestBuilder()).thenReturn(new RequestContext.Builder()); + when(sidecarClient.executeRequestAsync(any(RequestContext.class))) + .thenReturn(CompletableFuture.completedFuture(null)); + CreateSliceRequestPayload payload = mock(CreateSliceRequestPayload.class); + + api.createRestoreSliceFromExecutor(CLUSTER_ID, payload); + + ArgumentCaptor request = ArgumentCaptor.forClass(RequestContext.class); + verify(sidecarClient).executeRequestAsync(request.capture()); + assertThat(request.getValue().request().requestURI()).contains("\"MyKeyspace\"", "\"MyTable\"", JOB_ID.toString()); + assertThat(request.getValue().request().requestBody()).isSameAs(payload); + } + + @Test + void testRestoreJobProgressUsesQuotedIdentifiers() + { + RestoreJobProgressResponsePayload response = mock(RestoreJobProgressResponsePayload.class); + when(sidecarClient.restoreJobProgress(any(RestoreJobProgressRequestParams.class))) + .thenReturn(CompletableFuture.completedFuture(response)); + + api.restoreJobProgress(RestoreJobProgressFetchPolicy.FIRST_FAILED, cluster -> false, + (cluster, progress) -> { + assertThat(cluster).isEqualTo(CLUSTER_ID); + assertThat(progress).isSameAs(response); + }); + + ArgumentCaptor params = ArgumentCaptor.forClass(RestoreJobProgressRequestParams.class); + verify(sidecarClient).restoreJobProgress(params.capture()); + assertThat(params.getValue().keyspace).isEqualTo("\"MyKeyspace\""); + assertThat(params.getValue().table).isEqualTo("\"MyTable\""); + assertThat(params.getValue().jobId).isEqualTo(JOB_ID); + } +} diff --git a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/data/CassandraDataLayerTests.java b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/data/CassandraDataLayerTests.java index 088de9070..58f420b1e 100644 --- a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/data/CassandraDataLayerTests.java +++ b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/data/CassandraDataLayerTests.java @@ -21,13 +21,26 @@ import java.util.HashMap; import java.util.Map; +import java.util.concurrent.CompletableFuture; import com.google.common.collect.ImmutableMap; import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.CsvSource; +import o.a.c.sidecar.client.shaded.common.response.RingResponse; +import o.a.c.sidecar.client.shaded.common.response.TableStatsResponse; +import o.a.c.sidecar.client.shaded.common.response.data.RingEntry; +import o.a.c.sidecar.client.shaded.client.SidecarClient; +import o.a.c.sidecar.client.shaded.client.SidecarInstance; +import org.apache.cassandra.clients.Sidecar; + import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; class CassandraDataLayerTests { @@ -36,6 +49,35 @@ class CassandraDataLayerTests "table", "customers", "sidecar_contact_points", "localhost"); + @Test + void testDynamicSizingUsesQuotedIdentifiers() + { + Map options = new HashMap<>(REQUIRED_CLIENT_CONFIG_OPTIONS); + options.put("keyspace", "MyKeyspace"); + options.put("table", "MyTable"); + options.put(ClientConfig.SIZING_KEY, ClientConfig.SIZING_DYNAMIC); + options.put(ClientConfig.NUM_CORES_KEY, "2"); + options.put("consistencylevel", "ONE"); + ClientConfig clientConfig = ClientConfig.create(options); + CassandraDataLayer layer = new CassandraDataLayer(clientConfig, mock(Sidecar.ClientConfig.class), null); + layer.maybeQuotedKeyspace = "\"MyKeyspace\""; + layer.maybeQuotedTable = "\"MyTable\""; + layer.sidecar = mock(SidecarClient.class); + RingEntry entry = mock(RingEntry.class); + when(entry.fqdn()).thenReturn("localhost"); + RingResponse ring = new RingResponse(); + ring.add(entry); + TableStatsResponse stats = mock(TableStatsResponse.class); + when(stats.totalDiskSpaceUsedBytes()).thenReturn(1024L); + when(layer.sidecar.tableStats(any(SidecarInstance.class), eq("\"MyKeyspace\""), eq("\"MyTable\""))) + .thenReturn(CompletableFuture.completedFuture(stats)); + + assertThat(layer.getSizing(CompletableFuture.completedFuture(ring), + ReplicationFactor.simpleStrategy(1), clientConfig).getEffectiveNumberOfCores()) + .isEqualTo(1); + verify(layer.sidecar).tableStats(any(SidecarInstance.class), eq("\"MyKeyspace\""), eq("\"MyTable\"")); + } + @Test void testDefaultClearSnapshotStrategy() {