diff --git a/CHANGES.txt b/CHANGES.txt index 108096a416..1a06b877f4 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -10,6 +10,7 @@ * Add the ability to disable bulk loading of SSTables (CASSANDRA-18781) * Clean up obsolete functions and simplify cql_version handling in cqlsh (CASSANDRA-18787) Merged from 5.0: + * Fix storage_compatibility_mode for streaming (CASSANDRA-19126) * Add support of vector type to cqlsh COPY command (CASSANDRA-19118) * Make CQLSSTableWriter to support building of SAI indexes (CASSANDRA-18714) * Append additional JVM options when using JDK17+ (CASSANDRA-19001) diff --git a/build.xml b/build.xml index 79fdbda8d9..a4395d4ff5 100644 --- a/build.xml +++ b/build.xml @@ -1312,7 +1312,7 @@ - + @@ -1333,7 +1333,7 @@ - + diff --git a/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java b/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java index ed92ee6a41..5177fb09fb 100644 --- a/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java +++ b/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java @@ -316,14 +316,6 @@ public enum CassandraRelevantProperties /** Java Virtual Machine implementation name */ JAVA_VM_NAME("java.vm.name"), JOIN_RING("cassandra.join_ring", "true"), - /** - * {@link StorageCompatibilityMode} mode sets how the node will behave, sstable or messaging versions to use etc according to a yaml setting. - * But many tests don't load the config hence we need to force it otherwise they would run always under the default. Config is null for junits - * that don't load the config. Get from env var that CI/build.xml sets. - * - * This is a dev/CI only property. Do not use otherwise. - */ - JUNIT_STORAGE_COMPATIBILITY_MODE("cassandra.junit_storage_compatibility_mode", StorageCompatibilityMode.NONE.toString()), /** startup checks properties */ LIBJEMALLOC("cassandra.libjemalloc"), /** Line separator ("\n" on UNIX). */ @@ -588,6 +580,14 @@ public enum CassandraRelevantProperties TEST_SIMULATOR_PRINT_ASM_TYPES("cassandra.test.simulator.print_asm_types", ""), TEST_SKIP_CRYPTO_PROVIDER_INSTALLATION("cassandra.test.security.skip.provider.installation", "false"), TEST_SSTABLE_FORMAT_DEVELOPMENT("cassandra.test.sstableformatdevelopment"), + /** + * {@link StorageCompatibilityMode} mode sets how the node will behave, sstable or messaging versions to use etc according to a yaml setting. + * But many tests don't load the config hence we need to force it otherwise they would run always under the default. Config is null for junits + * that don't load the config. Get from env var that CI/build.xml sets. + * + * This is a dev/CI only property. Do not use otherwise. + */ + TEST_STORAGE_COMPATIBILITY_MODE("cassandra.test.storage_compatibility_mode", StorageCompatibilityMode.NONE.toString()), TEST_STRICT_LCS_CHECKS("cassandra.test.strict_lcs_checks"), /** Turns some warnings into exceptions for testing. */ TEST_STRICT_RUNTIME_CHECKS("cassandra.strict.runtime.checks"), diff --git a/src/java/org/apache/cassandra/config/Config.java b/src/java/org/apache/cassandra/config/Config.java index d62d836815..4db252917d 100644 --- a/src/java/org/apache/cassandra/config/Config.java +++ b/src/java/org/apache/cassandra/config/Config.java @@ -1280,7 +1280,7 @@ public class Config public double severity_during_decommission = 0; // TODO Revisit MessagingService::current_version - public StorageCompatibilityMode storage_compatibility_mode = StorageCompatibilityMode.NONE; + public StorageCompatibilityMode storage_compatibility_mode; /** * For the purposes of progress barrier we only support ALL, EACH_QUORUM, QUORUM, LOCAL_QUORUM, ANY, and ONE. diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index b3fbb283bb..f1440f1beb 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -239,6 +239,7 @@ public class DatabaseDescriptor private static ImmutableMap> sstableFormats; private static volatile SSTableFormat selectedSSTableFormat; + private static StorageCompatibilityMode storageCompatibilityMode = CassandraRelevantProperties.TEST_STORAGE_COMPATIBILITY_MODE.getEnum(true, StorageCompatibilityMode.class); private static Function commitLogSegmentMgrProvider = c -> DatabaseDescriptor.isCDCEnabled() ? new CommitLogSegmentManagerCDC(c, DatabaseDescriptor.getCommitLogLocation()) @@ -302,6 +303,8 @@ public class DatabaseDescriptor setConfig(loadConfig()); + applyCompatibilityMode(); + applySSTableFormats(); applySimpleConfig(); @@ -357,6 +360,7 @@ public class DatabaseDescriptor setDefaultFailureDetector(); Config.setClientMode(true); conf = configSupplier.get(); + applyCompatibilityMode(); diskOptimizationStrategy = new SpinningDiskOptimizationStrategy(); applySSTableFormats(); } @@ -445,7 +449,8 @@ public class DatabaseDescriptor private static void applyAll() throws ConfigurationException { - //InetAddressAndPort cares that applySimpleConfig runs first + applyCompatibilityMode(); + applySSTableFormats(); applyCryptoProvider(); @@ -1546,6 +1551,15 @@ public class DatabaseDescriptor return selectedFormat; } + private static void applyCompatibilityMode() + { + if (isClientInitialized()) + // tools or clients should not limit the sstable formats they support + storageCompatibilityMode = StorageCompatibilityMode.NONE; + else if (conf != null && conf.storage_compatibility_mode != null) + storageCompatibilityMode = conf.storage_compatibility_mode; + } + private static void applySSTableFormats() { ServiceLoader loader = ServiceLoader.load(SSTableFormat.Factory.class, DatabaseDescriptor.class.getClassLoader()); @@ -5021,11 +5035,7 @@ public class DatabaseDescriptor public static StorageCompatibilityMode getStorageCompatibilityMode() { - // Config is null for junits that don't load the config. Get from env var that CI/build.xml sets - if (conf == null) - return CassandraRelevantProperties.JUNIT_STORAGE_COMPATIBILITY_MODE.getEnum(StorageCompatibilityMode.NONE); - else - return conf.storage_compatibility_mode; + return storageCompatibilityMode; } public static ParameterizedClass getDefaultCompaction() diff --git a/src/java/org/apache/cassandra/net/MessagingService.java b/src/java/org/apache/cassandra/net/MessagingService.java index 8859d1a330..ceae703097 100644 --- a/src/java/org/apache/cassandra/net/MessagingService.java +++ b/src/java/org/apache/cassandra/net/MessagingService.java @@ -31,10 +31,10 @@ import java.util.stream.Collectors; import com.google.common.annotations.VisibleForTesting; import com.google.common.collect.Lists; -import io.netty.util.concurrent.Future; //checkstyle: permit this import import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import io.netty.util.concurrent.Future; //checkstyle: permit this import import org.apache.cassandra.concurrent.ScheduledExecutors; import org.apache.cassandra.concurrent.Stage; import org.apache.cassandra.config.DatabaseDescriptor; @@ -52,7 +52,6 @@ import org.apache.cassandra.utils.concurrent.Promise; import static java.util.Collections.synchronizedList; import static java.util.concurrent.TimeUnit.MINUTES; - import static org.apache.cassandra.concurrent.Stage.MUTATION; import static org.apache.cassandra.config.CassandraRelevantProperties.NON_GRACEFUL_SHUTDOWN; import static org.apache.cassandra.utils.Clock.Global.nanoTime; @@ -271,9 +270,26 @@ public class MessagingService extends MessagingServiceMBeanImpl implements Messa public static final int VERSION_50 = 13; // c14227 TTL overflow, 'uint' timestamps public static final int VERSION_51 = 14; // TCM public static final int minimum_version = VERSION_40; - public static final int current_version = Version.CURRENT.value; - static AcceptVersions accept_messaging = new AcceptVersions(minimum_version, current_version); - static AcceptVersions accept_streaming = new AcceptVersions(current_version, current_version); + public static final int maximum_version = VERSION_51; + // we want to use a modified behavior for the tools and clients - that is, since they are not running a server, they + // should not need to run in a compatibility mode. They should be able to connect to the server regardless whether + // it uses messaving version 4 or 5 + public static final int current_version = DatabaseDescriptor.getStorageCompatibilityMode().isBefore(5) ? VERSION_40 : VERSION_51; + static AcceptVersions accept_messaging; + static AcceptVersions accept_streaming; + static + { + if (DatabaseDescriptor.isClientInitialized()) + { + accept_messaging = new AcceptVersions(minimum_version, maximum_version); + accept_streaming = new AcceptVersions(minimum_version, maximum_version); + } + else + { + accept_messaging = new AcceptVersions(minimum_version, current_version); + accept_streaming = new AcceptVersions(current_version, current_version); + } + } static Map versionOrdinalMap = Arrays.stream(Version.values()).collect(Collectors.toMap(v -> v.value, v -> v.ordinal())); /** diff --git a/src/java/org/apache/cassandra/net/OutboundConnectionSettings.java b/src/java/org/apache/cassandra/net/OutboundConnectionSettings.java index 8a83a52aa6..ccc74f0aba 100644 --- a/src/java/org/apache/cassandra/net/OutboundConnectionSettings.java +++ b/src/java/org/apache/cassandra/net/OutboundConnectionSettings.java @@ -420,7 +420,7 @@ public class OutboundConnectionSettings if (tcpNoDelay != null) return tcpNoDelay; - if (isInLocalDC(getEndpointSnitch(), getBroadcastAddressAndPort(), to)) + if (DatabaseDescriptor.isClientOrToolInitialized() || isInLocalDC(getEndpointSnitch(), getBroadcastAddressAndPort(), to)) return INTRADC_TCP_NODELAY; return DatabaseDescriptor.getInterDCTcpNoDelay(); diff --git a/src/java/org/apache/cassandra/repair/messages/PrepareMessage.java b/src/java/org/apache/cassandra/repair/messages/PrepareMessage.java index a2951ef304..fc44d7d0c4 100644 --- a/src/java/org/apache/cassandra/repair/messages/PrepareMessage.java +++ b/src/java/org/apache/cassandra/repair/messages/PrepareMessage.java @@ -89,7 +89,9 @@ public class PrepareMessage extends RepairMessage } private static final String MIXED_MODE_ERROR = "Some nodes involved in repair are on an incompatible major version. " + - "Repair is not supported in mixed major version clusters (%d vs %d)."; + "Repair is not supported in mixed major version clusters (%d vs %d). Note that " + + "5.x nodes running in storage compatibility mode = 4 are considered " + + "4.x nodes."; public static final IVersionedSerializer serializer = new IVersionedSerializer() { diff --git a/src/java/org/apache/cassandra/tools/BulkLoader.java b/src/java/org/apache/cassandra/tools/BulkLoader.java index 367903a595..ecd6a5d31b 100644 --- a/src/java/org/apache/cassandra/tools/BulkLoader.java +++ b/src/java/org/apache/cassandra/tools/BulkLoader.java @@ -61,7 +61,7 @@ public class BulkLoader public static void load(LoaderOptions options) throws BulkLoadException { - DatabaseDescriptor.toolInitialization(); + DatabaseDescriptor.clientInitialization(); OutputHandler handler = new OutputHandler.SystemOutput(options.verbose, options.debug); SSTableLoader loader = new SSTableLoader( options.directory.toAbsolute(), diff --git a/test/unit/org/apache/cassandra/repair/FuzzTestBase.java b/test/unit/org/apache/cassandra/repair/FuzzTestBase.java index 6c676b6938..0c26b3f279 100644 --- a/test/unit/org/apache/cassandra/repair/FuzzTestBase.java +++ b/test/unit/org/apache/cassandra/repair/FuzzTestBase.java @@ -1201,7 +1201,7 @@ public abstract class FuzzTestBase extends CQLTester.InMemory { try (DataOutputBuffer b = DataOutputBuffer.scratchBuffer.get()) { - int messagingVersion = MessagingService.Version.CURRENT.value; + int messagingVersion = MessagingService.current_version; Message.serializer.serialize(msg, b, messagingVersion); DataInputBuffer in = new DataInputBuffer(b.unsafeGetBufferAndFlip(), false); return Message.serializer.deserialize(in, msg.from(), messagingVersion);