From 633bbee1a698b5666b459f7499fd100bd2ea3657 Mon Sep 17 00:00:00 2001 From: David Capwell Date: Thu, 4 Jan 2024 14:45:27 -0800 Subject: [PATCH] (Accord): Bug fixes from CASSANDRA-18675 to better support adding keyspaces patch by David Capwell; reviewed by Benedict Elliott Smith, Blake Eggleston for CASSANDRA-18804 --- modules/accord | 2 +- .../apache/cassandra/service/accord/AccordService.java | 7 +++++++ .../cassandra/service/accord/IAccordService.java | 9 --------- .../distributed/test/accord/AccordTestBase.java | 10 +++++++++- .../simulator/test/AccordJournalSimulationTest.java | 5 +++-- .../service/accord/AccordMessageSinkTest.java | 5 +++-- 6 files changed, 23 insertions(+), 15 deletions(-) diff --git a/modules/accord b/modules/accord index 0d8f60f742..901a0868cd 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit 0d8f60f742d443365a50115397ff1f0ab10fc694 +Subproject commit 901a0868cdaf6426226e6bafb0675773e04668bd diff --git a/src/java/org/apache/cassandra/service/accord/AccordService.java b/src/java/org/apache/cassandra/service/accord/AccordService.java index 93f5422f82..6cd1b68642 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordService.java +++ b/src/java/org/apache/cassandra/service/accord/AccordService.java @@ -28,6 +28,9 @@ import javax.annotation.Nonnull; import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Preconditions; + +import accord.coordinate.TopologyMismatch; +import org.apache.cassandra.cql3.statements.RequestValidations; import org.apache.cassandra.tcm.transformations.AddAccordTable; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -469,6 +472,10 @@ public class AccordService implements IAccordService, Shutdownable // Protocol also doesn't have a way to denote "unknown" outcome, so using a timeout as the closest match throw newPreempted(txnId, txn, consistencyLevel); } + if (cause instanceof TopologyMismatch) + { + throw RequestValidations.invalidRequest(cause.getMessage()); + } metrics.failures.mark(); throw new RuntimeException(cause); } diff --git a/src/java/org/apache/cassandra/service/accord/IAccordService.java b/src/java/org/apache/cassandra/service/accord/IAccordService.java index 2f0d7afc71..9422df4916 100644 --- a/src/java/org/apache/cassandra/service/accord/IAccordService.java +++ b/src/java/org/apache/cassandra/service/accord/IAccordService.java @@ -45,8 +45,6 @@ import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; import org.apache.cassandra.net.IVerbHandler; import org.apache.cassandra.net.Message; -import org.apache.cassandra.schema.Schema; -import org.apache.cassandra.schema.TableMetadata; import org.apache.cassandra.service.accord.api.AccordRoutableKey; import org.apache.cassandra.service.accord.api.AccordRoutingKey.TokenKey; import org.apache.cassandra.schema.TableId; @@ -156,13 +154,6 @@ public interface IAccordService void ensureTableIsAccordManaged(TableId tableId); - default void ensureTableIsAccordManaged(String keyspace, String table) - { - // TODO: remove when accord enabled is handled via schema - TableMetadata metadata = Schema.instance.getTableMetadata(keyspace, table); - ensureTableIsAccordManaged(metadata.id); - } - default void ensureKeyspaceIsAccordManaged(String keyspace) { // TODO: remove when accord enabled is handled via schema diff --git a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordTestBase.java b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordTestBase.java index 9cb218dab3..f5937f643b 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/accord/AccordTestBase.java +++ b/test/distributed/org/apache/cassandra/distributed/test/accord/AccordTestBase.java @@ -33,6 +33,8 @@ import java.util.stream.StreamSupport; import accord.coordinate.Invalidated; import com.google.common.base.Splitter; import com.google.common.primitives.Ints; +import org.apache.cassandra.schema.Schema; +import org.apache.cassandra.schema.TableMetadata; import org.junit.After; import org.junit.AfterClass; import org.junit.Before; @@ -134,7 +136,13 @@ public abstract class AccordTestBase extends TestBaseImpl public static void ensureTableIsAccordManaged(Cluster cluster, String ksname, String tableName) { - cluster.get(1).runOnInstance(() -> AccordService.instance().ensureTableIsAccordManaged(ksname, tableName)); + cluster.get(1).runOnInstance(() -> { + // TODO: remove when accord enabled is handled via schema + TableMetadata metadata = Schema.instance.getTableMetadata(ksname, tableName); + if (metadata == null) + return; // bad plumbing from shared utils.... + AccordService.instance().ensureTableIsAccordManaged(metadata.id); + }); } protected void test(List ddls, FailingConsumer fn) throws Exception diff --git a/test/simulator/test/org/apache/cassandra/simulator/test/AccordJournalSimulationTest.java b/test/simulator/test/org/apache/cassandra/simulator/test/AccordJournalSimulationTest.java index f4739bca6d..4201a13f7d 100644 --- a/test/simulator/test/org/apache/cassandra/simulator/test/AccordJournalSimulationTest.java +++ b/test/simulator/test/org/apache/cassandra/simulator/test/AccordJournalSimulationTest.java @@ -24,6 +24,8 @@ import java.util.concurrent.CopyOnWriteArrayList; import javax.annotation.Nullable; import com.google.common.collect.ImmutableMap; + +import accord.topology.TopologyUtils; import org.apache.cassandra.schema.*; import org.junit.Ignore; import org.junit.Test; @@ -35,7 +37,6 @@ import accord.api.Data; import accord.api.RoutingKey; import accord.api.Update; import accord.api.Write; -import accord.impl.TopologyUtils; import accord.local.Node; import accord.messages.PreAccept; import accord.messages.TxnRequest; @@ -208,7 +209,7 @@ public class AccordJournalSimulationTest extends SimulationTestBase { TxnId id = toTxnId(event); Ranges ranges = Ranges.of(new TokenRange(AccordRoutingKey.SentinelKey.min(tableId), AccordRoutingKey.SentinelKey.max(tableId))); - Topologies topologies = Utils.topologies(TopologyUtils.initialTopology(new Node.Id[] {node}, ranges, 3)); + Topologies topologies = Utils.topologies(TopologyUtils.initialTopology(new Node.Id[] { node}, ranges, 3)); Keys keys = Keys.of(toKey(0)); Txn txn = new Txn.InMemory(keys, new TxnRead(new TxnNamedRead[0], keys, null), TxnQuery.ALL, new NoopUpdate()); FullRoute route = route(); diff --git a/test/unit/org/apache/cassandra/service/accord/AccordMessageSinkTest.java b/test/unit/org/apache/cassandra/service/accord/AccordMessageSinkTest.java index a35050b017..93150d1054 100644 --- a/test/unit/org/apache/cassandra/service/accord/AccordMessageSinkTest.java +++ b/test/unit/org/apache/cassandra/service/accord/AccordMessageSinkTest.java @@ -20,6 +20,8 @@ package org.apache.cassandra.service.accord; import org.junit.BeforeClass; import org.junit.Test; + +import accord.topology.TopologyUtils; import org.mockito.ArgumentCaptor; import org.mockito.Mockito; @@ -27,7 +29,6 @@ import accord.Utils; import accord.api.Agent; import accord.impl.AbstractFetchCoordinator; import accord.impl.IntKey; -import accord.impl.TopologyUtils; import accord.local.Node; import accord.messages.InformOfTxnId; import accord.messages.MessageType; @@ -56,7 +57,7 @@ public class AccordMessageSinkTest { private static final Node.Id node = new Node.Id(1); private static final AccordEndpointMapper mapping = SimpleAccordEndpointMapper.INSTANCE; - private static final Topology topology = TopologyUtils.initialTopology(new Node.Id[] {node}, Ranges.of(IntKey.range(0, 100)), 1); + private static final Topology topology = TopologyUtils.initialTopology(new Node.Id[] { node}, Ranges.of(IntKey.range(0, 100)), 1); private static final Topologies topologies = new Topologies.Single((a, b, ignore) -> 0, topology); private static final MessageDelivery messaging = Mockito.mock(MessageDelivery.class);