diff --git a/CHANGES.txt b/CHANGES.txt index 1b24ae81c5..b458312d7c 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -2,6 +2,7 @@ * Allow nodetool garbagecollect to take a user defined list of SSTables (CASSANDRA-16767) * Add a guardrail for misprepared statements (CASSANDRA-21139) Merged from 6.0: + * Add an offline cluster metadata tool (CASSANDRA-19151) * Ensure schema created before 2.1 without tableId in folder name can be loaded in SnapshotLoader (CASSANDRA-21246) * Differentiate between legitimate cases where the first entry is the same as the last entry and empty bounds in SSTableCursorWriter#addIndexBlock() (CASSANDRA-21255) * Introduce minimum_threshold for data resurrection startup check (CASSANDRA-21293) diff --git a/src/java/org/apache/cassandra/tools/CMSOfflineTool.java b/src/java/org/apache/cassandra/tools/CMSOfflineTool.java new file mode 100644 index 0000000000..90d44e3fc0 --- /dev/null +++ b/src/java/org/apache/cassandra/tools/CMSOfflineTool.java @@ -0,0 +1,825 @@ +/* + * 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.tools; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.Collection; +import java.util.EnumSet; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import java.util.stream.Collectors; + +import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.dht.IPartitioner; +import org.apache.cassandra.dht.Range; +import org.apache.cassandra.dht.Token; +import org.apache.cassandra.io.util.File; +import org.apache.cassandra.io.util.FileInputStreamPlus; +import org.apache.cassandra.io.util.FileOutputStreamPlus; +import org.apache.cassandra.locator.InetAddressAndPort; +import org.apache.cassandra.locator.MetaStrategy; +import org.apache.cassandra.locator.Replica; +import org.apache.cassandra.schema.KeyspaceMetadata; +import org.apache.cassandra.schema.ReplicationParams; +import org.apache.cassandra.tcm.ClusterMetadata; +import org.apache.cassandra.tcm.ClusterMetadataService; +import org.apache.cassandra.tcm.MultiStepOperation; +import org.apache.cassandra.tcm.membership.Directory; +import org.apache.cassandra.tcm.membership.Location; +import org.apache.cassandra.tcm.membership.NodeAddresses; +import org.apache.cassandra.tcm.membership.NodeId; +import org.apache.cassandra.tcm.membership.NodeState; +import org.apache.cassandra.tcm.membership.NodeVersion; +import org.apache.cassandra.tcm.ownership.DataPlacement; +import org.apache.cassandra.tcm.ownership.ReplicaGroups; +import org.apache.cassandra.tcm.ownership.UniformRangePlacement; +import org.apache.cassandra.tcm.sequences.BootstrapAndJoin; +import org.apache.cassandra.tcm.sequences.Move; +import org.apache.cassandra.tcm.sequences.ReconfigureCMS; +import org.apache.cassandra.tcm.serialization.VerboseMetadataSerializer; +import org.apache.cassandra.tcm.serialization.Version; +import org.apache.cassandra.tcm.transformations.Assassinate; +import org.apache.cassandra.tcm.transformations.CancelInProgressSequence; +import org.apache.cassandra.tcm.transformations.PrepareMove; +import org.apache.cassandra.tcm.transformations.Unregister; +import org.apache.cassandra.tcm.transformations.UnsafeJoin; +import org.apache.cassandra.tcm.transformations.cms.PrepareCMSReconfiguration; +import org.apache.cassandra.utils.FBUtilities; + +import picocli.CommandLine; +import picocli.CommandLine.Command; +import picocli.CommandLine.Option; + +import static com.google.common.base.Throwables.getStackTraceAsString; +import static org.apache.cassandra.tcm.transformations.cms.PrepareCMSReconfiguration.needsReconfiguration; + +/** + * Offline tool to print or update cluster metadata stored in a dump file. + *

+ * The tool operates entirely offline: it reads a metadata dump file produced by + * {@code nodetool cms dump} (or equivalent), applies the requested transformation, + * and writes the result to a new file. The original dump file is never modified. + *

+ * Run without a subcommand to print usage information. + */ +@SuppressWarnings({ "unused", "DefaultAnnotationParam", "UseOfSystemOutOrSystemErr" }) +@Command(name = "cmsofflinetool", +mixinStandardHelpOptions = true, +description = "Offline tool to print or update cluster metadata dump.", +subcommands = { CMSOfflineTool.AbortBootstrap.class, + CMSOfflineTool.AbortDecommission.class, + CMSOfflineTool.AbortMove.class, + CMSOfflineTool.AssassinateNode.class, + CMSOfflineTool.Describe.class, + CMSOfflineTool.ForceJoin.class, + CommandLine.HelpCommand.class, + CMSOfflineTool.MoveToken.class, + CMSOfflineTool.Print.class, + CMSOfflineTool.PrintDataPlacements.class, + CMSOfflineTool.PrintDirectoryCmd.class, + CMSOfflineTool.ResetCMS.class }) +public class CMSOfflineTool implements Runnable +{ + private final Output output; + + public CMSOfflineTool(Output output) + { + this.output = output; + } + + public static void main(String[] args) throws IOException + { + CMSOfflineTool tool = new CMSOfflineTool(new Output(System.out, System.err)); + CommandLine cli = new CommandLine(tool) + .setColorScheme(CommandLine.Help.defaultColorScheme(CommandLine.Help.Ansi.OFF)) + .setExecutionExceptionHandler((ex, cmd, parseResult) -> { + cmd.getErr().println("Error: " + ex.getMessage()); + cmd.getErr().println("-- StackTrace --"); + cmd.getErr().println(getStackTraceAsString(ex)); + return 2; + }); + int status = cli.execute(args); + System.exit(status); + } + + @Override + public void run() + { + CommandLine.usage(this, output.out, CommandLine.Help.Ansi.OFF); + } + + public static abstract class ClusterMetadataToolCmd implements Runnable + { + @Option(names = { "-f", "--file" }, description = "Cluster metadata dump file path.", required = true) + protected String metadataDumpFile; + + @Option(names = { "-sv", "--serialization-version" }, description = "Serialization version to use.") + private Version serializationVersion; + + @CommandLine.ParentCommand + private CMSOfflineTool parent; + + @Override + public void run() + { + try + { + execute(parent.output); + parent.output.out.flush(); + parent.output.err.flush(); + } + catch (IOException e) + { + throw new RuntimeException(e); + } + } + + protected abstract void execute(Output output) throws IOException; + + public ClusterMetadata parseClusterMetadata() throws IOException + { + File file = new File(metadataDumpFile); + if (!file.exists()) + { + throw new IllegalArgumentException("Cluster metadata dump file " + metadataDumpFile + " does not exist."); + } + + // Make sure the partitioner we use to manipulate the metadata is the same one used to generate it + IPartitioner partitioner; + try (FileInputStreamPlus fisp = new FileInputStreamPlus(metadataDumpFile)) + { + int x = fisp.readUnsignedVInt32(); + Version version = Version.fromInt(x); + partitioner = ClusterMetadata.Serializer.getPartitioner(fisp, version); + } + DatabaseDescriptor.toolInitialization(); + DatabaseDescriptor.setPartitionerUnsafe(partitioner); + ClusterMetadataService.initializeForTools(false); + + return ClusterMetadataService.deserializeClusterMetadata(metadataDumpFile); + } + + public void writeMetadata(Output output, ClusterMetadata metadata, String outputFilePath) throws IOException + { + Version serializationVersion = getSerializationVersion(metadata); + Path p = outputFilePath != null + ? Files.createFile(Path.of(outputFilePath)) + : Files.createTempFile("clustermetadata", ".dump"); + + try (FileOutputStreamPlus out = new FileOutputStreamPlus(p)) + { + VerboseMetadataSerializer.serialize(ClusterMetadata.serializer, metadata, out, serializationVersion); + output.out.println("Updated cluster metadata written to file " + p.toAbsolutePath()); + } + } + + Version getSerializationVersion(ClusterMetadata metadata) + { + // Step 1: Default to the serialization version from the cluster metadata + Version finalVersion = metadata.directory.commonSerializationVersion; + + // Step 2: User-specified version takes precedence + if (serializationVersion != null) + { + // Warn if user-specified version is older than what the metadata was written with + if (serializationVersion.isBefore(finalVersion)) + { + parent.output.err.printf("WARNING: Given serialization version %s is older than " + + "the version in cluster metadata (%s). Proceeding as requested.%n", + serializationVersion, finalVersion); + } + finalVersion = serializationVersion; + } + + // Step 3: Current binary version must be able to handle the finalized version + Version currentVersion = NodeVersion.CURRENT.serializationVersion(); + if (currentVersion.isBefore(finalVersion)) + { + throw new IllegalArgumentException("Current version " + currentVersion + + " is older than the target serialization version " + + finalVersion + ". Cannot proceed further. " + + "Try modifying cluster metadata using binaries that support " + + "minimum serialization version: " + finalVersion + '.'); + } + + return finalVersion; + } + } + + /** + * Base class for commands that cancel an in-progress sequence for a given node. + * Subclasses specify which sequence kinds they handle via {@link #supportedKinds()} and + * may apply additional transformations after cancellation via {@link #postCancel}. + */ + abstract static class AbstractAbortSequence extends ClusterMetadataToolCmd + { + @CommandLine.ArgGroup(exclusive = true, multiplicity = "1") + NodeIdentifierOption nodeIdentifierOption; + + @Option(names = { "-o", "--output-file" }, + description = "Output file path for storing the updated Cluster Metadata.") + String outputFilePath; + + protected abstract EnumSet supportedKinds(); + + protected ClusterMetadata postCancel(ClusterMetadata metadata, NodeId nodeId) + { + return metadata; + } + + @Override + protected void execute(Output output) throws IOException + { + ClusterMetadata metadata = parseClusterMetadata(); + NodeId nodeId = nodeIdentifierOption.getNodeId(metadata); + + MultiStepOperation multiStepOperation = metadata.inProgressSequences.get(nodeId); + if (multiStepOperation == null) + { + throw new IllegalArgumentException("No transformation sequence is in progress for " + + nodeIdentifierOption.getNodeIpOrId() + '.'); + } + + if (!supportedKinds().contains(multiStepOperation.kind())) + { + throw new IllegalArgumentException("Sequence of kind " + multiStepOperation.kind() + + " is in progress for node " + nodeIdentifierOption.getNodeIpOrId() + + ". Cannot proceed with this operation."); + } + + CancelInProgressSequence cancelSequence = new CancelInProgressSequence(nodeId); + ClusterMetadata updatedMetadata = cancelSequence.execute(metadata).success().metadata; + updatedMetadata = postCancel(updatedMetadata, nodeId); + writeMetadata(output, updatedMetadata, outputFilePath); + } + } + + /** + * Cancels a JOIN or REPLACE bootstrap sequence for the given node and unregisters it + * from the cluster. Use this when a node is stuck in bootstrapping or replacement. + * Fails if no in-progress sequence exists, or if the sequence is not of kind JOIN or REPLACE. + */ + @Command(name = "abortbootstrap", + description = "Aborts bootstrap for given node if in progress.") + public static class AbortBootstrap extends AbstractAbortSequence + { + @Override + protected EnumSet supportedKinds() + { + return EnumSet.of(MultiStepOperation.Kind.JOIN, MultiStepOperation.Kind.REPLACE); + } + + @Override + protected ClusterMetadata postCancel(ClusterMetadata metadata, NodeId nodeId) + { + // Cancelling the sequence is not enough, we need to unregister as well + Unregister unregister = new Unregister(nodeId, EnumSet.of(NodeState.REGISTERED), + ClusterMetadataService.instance().placementProvider()); + return unregister.execute(metadata).success().metadata; + } + } + + /** + * Cancels an in-progress MOVE sequence for the given node, returning it to its + * pre-move token assignment. Fails if no in-progress sequence exists, or if the + * sequence is not of kind MOVE. + */ + @Command(name = "abortmove", description = "Aborts in progress move sequence for given node.") + static class AbortMove extends AbstractAbortSequence + { + @Override + protected EnumSet supportedKinds() + { + return EnumSet.of(MultiStepOperation.Kind.MOVE); + } + } + + /** + * Cancels an in-progress LEAVE (decommission) sequence for the given node, keeping + * it as an active member of the ring. Fails if no in-progress sequence exists, or if + * the sequence is not of kind LEAVE. + */ + @Command(name = "abortdecommission", description = "Aborts in progress decommission sequence for given node.") + static class AbortDecommission extends AbstractAbortSequence + { + @Override + protected EnumSet supportedKinds() + { + return EnumSet.of(MultiStepOperation.Kind.LEAVE); + } + } + + /** + * Removes a node from cluster metadata by applying the {@link Assassinate} transformation. + * If a MOVE or LEAVE sequence is in progress for the node, it is cancelled first. + * If the node is a CMS member, it is removed from CMS before assassination. + * Fails if the node is in a JOIN or REPLACE sequence; use {@code abortbootstrap} instead. + */ + @Command(name = "assassinate", description = "Assassinates given node from Cluster metadata.") + static class AssassinateNode extends ClusterMetadataToolCmd + { + private final EnumSet supportedCancelSequences = + EnumSet.of(MultiStepOperation.Kind.MOVE, MultiStepOperation.Kind.LEAVE); + @CommandLine.ArgGroup(exclusive = true, multiplicity = "1") + NodeIdentifierOption nodeIdentifierOption; + @Option(names = { "-o", "--output-file" }, + description = "Output file path for storing the updated Cluster Metadata.") + private String outputFilePath; + + @Override + protected void execute(Output output) throws IOException + { + ClusterMetadata metadata = parseClusterMetadata(); + NodeId nodeId = nodeIdentifierOption.getNodeId(metadata); + + // Check if there are any in-progress sequences for given node + // If any, then cancel the sequence and then assassinate it + if (metadata.inProgressSequences.contains(nodeId)) + { + MultiStepOperation multiStepOperation = metadata.inProgressSequences.get(nodeId); + MultiStepOperation.Kind sequenceKind = multiStepOperation.kind(); + if (!supportedCancelSequences.contains(sequenceKind)) + { + if (sequenceKind == MultiStepOperation.Kind.JOIN || sequenceKind == MultiStepOperation.Kind.REPLACE) + { + throw new IllegalArgumentException("Cannot assassinate the node when sequence of kind " + + sequenceKind + " is in progress. " + + "Run abortbootstrap instead."); + } + else + { + throw new IllegalArgumentException("Cannot assassinate the node when sequence of kind " + + sequenceKind + " is in progress."); + } + } + + output.out.printf("Cancelling in-progress sequence of kind %s before assassinating node.\n", + metadata.inProgressSequences.get(nodeId).kind()); + metadata = new CancelInProgressSequence(nodeId).execute(metadata).success().metadata; + } + + metadata = maybeReconfigureCMS(metadata, nodeId); + Assassinate transformation = new Assassinate(nodeId, ClusterMetadataService.instance().placementProvider()); + ClusterMetadata updatedMetadata = transformation.execute(metadata).success().metadata; + + writeMetadata(output, updatedMetadata, outputFilePath); + } + + ClusterMetadata maybeReconfigureCMS(ClusterMetadata metadata, NodeId nodeId) + { + InetAddressAndPort addressAndPort = metadata.directory.endpoint(nodeId); + if (!metadata.isCMSMember(addressAndPort)) + { + return metadata; + } + + // Ref: org.apache.cassandra.tcm.sequences.ReconfigureCMS.maybeReconfigureCMS + PrepareCMSReconfiguration.Simple transformation = new PrepareCMSReconfiguration.Simple(nodeId, Set.of()); + ClusterMetadata updatedMetadata = transformation.execute(metadata).success().metadata; + updatedMetadata = updatedMetadata.inProgressSequences.get(ReconfigureCMS.SequenceKey.instance) + .applyTo(updatedMetadata).success().metadata; + if (updatedMetadata.isCMSMember(addressAndPort)) + { + throw new IllegalStateException("Could not remove node " + nodeIdentifierOption.getNodeIpOrId() + + " from CMS."); + } + return updatedMetadata; + } + } + + /** + * Replaces the entire CMS membership with a single node. All existing CMS replicas are + * removed and the specified node becomes the sole CMS member. Useful for recovering from + * a state where the CMS is unreachable. + */ + @Command(name = "resetcms", + description = "Replaces all CMS members with the specified node. WARNING: all existing CMS replicas are removed.") + static class ResetCMS extends ClusterMetadataToolCmd + { + @CommandLine.ArgGroup(exclusive = true, multiplicity = "1") + NodeIdentifierOption nodeIdentifierOption; + + @Option(names = { "-o", "--output-file" }, + description = "Output file path for storing the updated Cluster Metadata.") + private String outputFilePath; + + @Override + protected void execute(Output output) throws IOException + { + ClusterMetadata metadata = parseClusterMetadata(); + NodeId nodeId = nodeIdentifierOption.getNodeId(metadata); + metadata = resetCMS(metadata, nodeId); + writeMetadata(output, metadata, outputFilePath); + } + + ClusterMetadata resetCMS(ClusterMetadata metadata, NodeId nodeId) + { + NodeState nodeState = metadata.directory.peerState(nodeId); + if (nodeState != NodeState.JOINED) + { + throw new IllegalArgumentException("Node " + nodeIdentifierOption.getNodeIpOrId() + " is in " + + nodeState + " state. Only a JOINED node can be set as CMS member."); + } + InetAddressAndPort endpoint = metadata.directory.getNodeAddresses(nodeId).broadcastAddress; + ReplicationParams metaParams = ReplicationParams.meta(metadata); + Iterable currentReplicas = metadata.placements.get(metaParams).writes.byEndpoint().flattenValues(); + DataPlacement.Builder placementBuilder = metadata.placements.get(metaParams).unbuild(); + for (Replica replica : currentReplicas) + { + placementBuilder.withoutReadReplica(metadata.epoch, replica) + .withoutWriteReplica(metadata.epoch, replica); + } + + Replica newCMS = MetaStrategy.replica(endpoint); + placementBuilder.withReadReplica(metadata.epoch, newCMS) + .withWriteReplica(metadata.epoch, newCMS); + + return metadata.transformer() + .with(metadata.placements.unbuild() + .with(metaParams, placementBuilder.build()) + .build()) + .build().metadata; + } + } + + /** + * Moves a node to a new token. Only supports single-token nodes. + * If a MOVE sequence is already in progress for the node, it is completed rather than started fresh; + * in that case the provided token must match the token of the in-progress sequence. + */ + @Command(name = "move", + description = "Moves node to given token. Works only for cluster having nodes with single token.") + static class MoveToken extends ClusterMetadataToolCmd + { + @CommandLine.ArgGroup(exclusive = true, multiplicity = "1") + NodeIdentifierOption nodeIdentifierOption; + + @Option(names = { "-t", "--token" }, + description = "Token to assign.") + private String token; + + @Option(names = { "-o", "--output-file" }, + description = "Output file path for storing the updated Cluster Metadata.") + private String outputFilePath; + + @Override + protected void execute(Output output) throws IOException + { + // Took the reference from org.apache.cassandra.tcm.sequences.SingleNodeSequences.move + ClusterMetadata metadata = parseClusterMetadata(); + NodeId nodeId = nodeIdentifierOption.getNodeId(metadata); + + if (metadata.inProgressSequences.contains(nodeId)) + { + ClusterMetadata updatedMetadata = finishInProgressSequence(nodeId, token, metadata); + writeMetadata(output, updatedMetadata, outputFilePath); + return; + } + + if (null == token) + { + throw new IllegalArgumentException("Token required when no MOVE sequence is in progress."); + } + metadata.partitioner.getTokenFactory().validate(token); + Token toToken = metadata.partitioner.getTokenFactory().fromString(token); + if (metadata.tokenMap.tokens().contains(toToken)) + { + NodeId tokenOwnerId = metadata.tokenMap.owner(toToken); + throw new IllegalArgumentException("Target token " + toToken + " is already owned by node " + tokenOwnerId.id()); + } + if (metadata.tokenMap.tokens(nodeId).size() > 1) + { + throw new UnsupportedOperationException("This node has more than one token and cannot be moved thusly."); + } + + PrepareMove prepareMove = new PrepareMove(nodeId, Set.of(toToken), new UniformRangePlacement(), false); + ClusterMetadata updatedMetadata = prepareMove.execute(metadata).success().metadata; + updatedMetadata = updatedMetadata.inProgressSequences.get(nodeId).applyTo(updatedMetadata).success().metadata; + writeMetadata(output, updatedMetadata, outputFilePath); + } + + ClusterMetadata finishInProgressSequence(NodeId nodeId, String token, ClusterMetadata metadata) + { + MultiStepOperation multiStepOperation = metadata.inProgressSequences.get(nodeId); + if (multiStepOperation.kind() != MultiStepOperation.Kind.MOVE) + { + throw new IllegalArgumentException("Another sequence of kind " + multiStepOperation.kind() + + " is in progress for node " + nodeIdentifierOption.getNodeIpOrId() + + ". Cannot proceed with move."); + } + + Move moveInProgress = (Move) multiStepOperation; + Collection inProgressTokens = moveInProgress.tokens; + Collection givenTokenSet; + if (token == null) + { + givenTokenSet = Set.of(); + } + else + { + metadata.partitioner.getTokenFactory().validate(token); + givenTokenSet = Set.of(metadata.partitioner.getTokenFactory().fromString(token)); + } + if (givenTokenSet.isEmpty() || new HashSet<>(moveInProgress.tokens).equals(givenTokenSet)) + { + return moveInProgress.applyTo(metadata).success().metadata; + } + + throw new IllegalArgumentException("Move in progress for another token(s) " + inProgressTokens + '.'); + } + } + + /** + * Identifies a target node using either its IP address or its integer/UUID node ID. + * Exactly one of {@code -ip} or {@code -id} must be provided. + */ + static class NodeIdentifierOption + { + @Option(names = { "-ip", "--ip-address" }, required = true, + description = "IP address of the target endpoint. Port can be optionally specified " + + "using a colon after the IP address (e.g., 127.0.0.1:9042).") + private String ip; + + @Option(names = { "-id", "--node-id" }, required = true, + description = "Node ID. It can be integer ID assigned to node or the node uuid.") + private String id; + + String getNodeIpOrId() + { + return ip != null ? ip : id; + } + + NodeId getNodeId(ClusterMetadata metadata) + { + if (id != null) + { + NodeId nodeId = NodeId.fromString(id); + if (!metadata.directory.peerIds().contains(nodeId)) + { + throw new IllegalArgumentException("No node present with id " + id + + " in the given cluster metadata."); + } + return nodeId; + } + else if (ip != null) + { + InetAddressAndPort nodeAddress = InetAddressAndPort.getByNameUnchecked(ip); + NodeId nodeId = metadata.directory.peerId(nodeAddress); + + if (null == nodeId) + { + throw new IllegalArgumentException("No node present with ip address " + ip + + " in given cluster metadata."); + } + return nodeId; + } + + throw new IllegalArgumentException("Neither node id nor ip address specified to fetch NodeId from metadata."); + } + } + + + /** + * Prints a high-level summary of the CMS state: members, epoch, replication factor, + * service state, and whether reconfiguration is needed. + */ + @Command(name = "describe", description = "Describes the cluster metadata.") + static class Describe extends ClusterMetadataToolCmd + { + @Override + protected void execute(Output output) throws IOException + { + ClusterMetadata metadata = parseClusterMetadata(); + String members = metadata.fullCMSMembers() + .stream() + .sorted() + .map(Object::toString) + .collect(Collectors.joining(",")); + + output.out.printf("Cluster Metadata Service:%n"); + output.out.printf("Members: %s%n", members); + output.out.printf("Needs reconfiguration: %s%n", needsReconfiguration(metadata)); + output.out.printf("Service State: %s%n", ClusterMetadataService.state(metadata)); + output.out.printf("Epoch: %s%n", metadata.epoch.getEpoch()); + output.out.printf("Replication factor: %s%n", ReplicationParams.meta(metadata).toString()); + } + } + + /** + * Forces a node directly to JOINED state, bypassing the normal bootstrap sequence. + *

+ * Fails if the node is already JOINED, or if a non-JOIN sequence is in progress. + */ + @Command(name = "forcejoin", description = "Forces a node to move to JOINED state.") + static class ForceJoin extends ClusterMetadataToolCmd + { + @SuppressWarnings("MismatchedQueryAndUpdateOfCollection") + @Option(names = { "-t", "--token" }, + description = "Token to assign. Pass it multiple times to assign multiple tokens to node.") + private final List tokens = new ArrayList<>(); + @CommandLine.ArgGroup(exclusive = true, multiplicity = "1") + NodeIdentifierOption nodeIdentifierOption; + @Option(names = { "-o", "--output-file" }, + description = "Output file path for storing the updated Cluster Metadata.") + private String outputFilePath; + + @Override + protected void execute(Output output) throws IOException + { + ClusterMetadata metadata = parseClusterMetadata(); + NodeId nodeId = nodeIdentifierOption.getNodeId(metadata); + Set tokenSet = new HashSet<>(tokens.size()); + Token.TokenFactory tokenFactory = metadata.partitioner.getTokenFactory(); + tokens.forEach(t -> tokenSet.add(tokenFactory.fromString(t))); + + NodeState nodeState = metadata.directory.peerState(nodeId); + if (nodeState == NodeState.JOINED) + { + throw new IllegalArgumentException("Node " + nodeIdentifierOption.getNodeIpOrId() + + " is already in JOINED state."); + } + + ClusterMetadata updatedMetadata; + if (metadata.inProgressSequences.get(nodeId) != null) + { + MultiStepOperation multiStepOperation = metadata.inProgressSequences.get(nodeId); + if (multiStepOperation.kind() != MultiStepOperation.Kind.JOIN) + { + throw new IllegalArgumentException("Another sequence of kind " + multiStepOperation.kind() + + " is in progress for node " + nodeIdentifierOption.getNodeIpOrId() + + ". Cannot proceed with force join."); + } + BootstrapAndJoin bootstrapAndJoin = (BootstrapAndJoin) multiStepOperation; + Set sequenceTokens = bootstrapAndJoin.finishJoin.tokens; + if (tokens.isEmpty() + || (tokenSet.size() == sequenceTokens.size() && sequenceTokens.containsAll(tokenSet))) + { + updatedMetadata = bootstrapAndJoin.applyTo(metadata).success().metadata; + } + else + { + // If tokens are provided, then it should match with the in-progress sequence tokens + throw new IllegalArgumentException("The tokens provided " + tokens + " do not match with " + + " in progress BootstrapAndJoin sequence tokens " + + sequenceTokens + ". Cannot proceed further."); + } + } + else + { + // There are no in-progress sequences, force join by using UnsafeJoin transformation + if (tokenSet.isEmpty()) + { + throw new IllegalArgumentException("Tokens must be provided to force join a node."); + } + UnsafeJoin unsafeJoin = new UnsafeJoin(nodeId, tokenSet, new UniformRangePlacement()); + updatedMetadata = unsafeJoin.execute(metadata).success().metadata; + } + + writeMetadata(output, updatedMetadata, outputFilePath); + } + } + + /** + * Prints the full {@code toString()} representation of the cluster metadata to stdout. + * Useful for a quick human-readable overview of the entire metadata state. + *

+ * Note: The output format is not stable and may change between versions. + * Do not rely on it for programmatic parsing. + */ + @Command(name = "print", description = "Prints string output of the cluster metadata. " + + "Output format is subject to change and should not be relied on for parsing.") + static class Print extends ClusterMetadataToolCmd + { + @Override + protected void execute(Output output) throws IOException + { + // It supports only toString output of the ClusterMetadata for now + ClusterMetadata metadata = parseClusterMetadata(); + output.out.println(metadata); + } + } + + /** + * Prints per-node directory information: addresses, ports, rack, DC, state, + * serialization version, and CMS membership status. + */ + @Command(name = "printdirectory", description = "Prints directory information in cluster metadata file.") + static class PrintDirectoryCmd extends ClusterMetadataToolCmd + { + @Override + protected void execute(Output output) throws IOException + { + ClusterMetadata metadata = parseClusterMetadata(); + Directory directory = metadata.directory; + Set nodeIdList = directory.peerIds(); + for (NodeId nodeId : nodeIdList) + { + NodeAddresses nodeAddresses = directory.getNodeAddresses(nodeId); + Location location = directory.location(nodeId); + output.out.println("NodeId: " + nodeId.id()); + String format = " %-22s%s\n"; + output.out.printf(format, "rack", location.rack); + output.out.printf(format, "local_port", nodeAddresses.localAddress.getPort()); + output.out.printf(format, "broadcast_port", nodeAddresses.broadcastAddress.getPort()); + output.out.printf(format, "host_id", nodeId.toUUID()); + output.out.printf(format, "broadcast_address", nodeAddresses.broadcastAddress.getAddress().toString()); + output.out.printf(format, "native_address", nodeAddresses.nativeAddress.getAddress().toString()); + output.out.printf(format, "native_port", nodeAddresses.nativeAddress.getPort()); + output.out.printf(format, "local_address", nodeAddresses.localAddress.getAddress().toString()); + output.out.printf(format, "state", directory.peerState(nodeId)); + output.out.printf(format, "serialization_version", directory.version(nodeId).serializationVersion()); + output.out.printf(format, "cassandra_version", directory.version(nodeId).cassandraVersion); + output.out.printf(format, "dc", location.datacenter); + output.out.printf(format, "is_cms_member", metadata.isCMSMember(nodeAddresses.broadcastAddress)); + } + } + } + + /** + * Prints read and write replica placements for a specific keyspace, sorted by token range. + * Requires the {@code -ks} option to specify the target keyspace. + */ + @Command(name = "printdataplacements", description = "Prints data placements in cluster metadata file.") + static class PrintDataPlacements extends ClusterMetadataToolCmd + { + @Option(names = { "-ks", "--keyspace" }, required = true, + description = "Keyspace to use for printing data placements.") + private String keyspace; + + @Override + protected void execute(Output output) throws IOException + { + ClusterMetadata metadata = parseClusterMetadata(); + + KeyspaceMetadata keyspaceMetadata = metadata.schema.getKeyspaces().getNullable(keyspace); + if (keyspaceMetadata == null) + { + throw new IllegalArgumentException("Keyspace " + keyspace + " not found in cluster metadata."); + } + + DataPlacement placement = metadata.placements.get(keyspaceMetadata.params.replication); + List rows = new ArrayList<>(); + rows.addAll(replicaGroupsToRows(placement.reads, "read")); + rows.addAll(replicaGroupsToRows(placement.writes, "write")); + + rows.sort((o1, o2) -> { + Range left = (Range) o1[0]; + Range right = (Range) o2[0]; + return left.compareTo(right); + }); + + int rangeMaxLength = 0; + for (Object[] objects : rows) + { + rangeMaxLength = Math.max(rangeMaxLength, objects[0].toString().length()); + } + + String rowFormat = String.format("%%-%ds %%-7s %%s\n", (rangeMaxLength + 2)); + output.out.printf(rowFormat, "Token Range", "Type", "Endpoints"); + + rows.forEach(row -> output.out.printf(rowFormat, row)); + } + + List replicaGroupsToRows(ReplicaGroups replicaGroups, String replicaGroupType) + { + List rows = new ArrayList<>(); + replicaGroups.forEach(((tokenRange, forRange) -> { + List endpoints = new ArrayList<>(forRange.get().size()); + forRange.get().forEach(replica -> endpoints.add(replica.endpoint().toString())); + String addresses = String.join(", ", endpoints); + rows.add(new Object[]{ tokenRange, replicaGroupType, addresses }); + })); + return rows; + } + } + + static + { + FBUtilities.preventIllegalAccessWarnings(); + } +} diff --git a/src/java/org/apache/cassandra/tools/nodetool/CMSAdmin.java b/src/java/org/apache/cassandra/tools/nodetool/CMSAdmin.java index 20aad56a17..364e5873d4 100644 --- a/src/java/org/apache/cassandra/tools/nodetool/CMSAdmin.java +++ b/src/java/org/apache/cassandra/tools/nodetool/CMSAdmin.java @@ -18,6 +18,7 @@ package org.apache.cassandra.tools.nodetool; +import java.io.IOException; import java.io.PrintStream; import java.util.ArrayList; import java.util.Comparator; @@ -28,9 +29,11 @@ import java.util.Map; import com.google.common.collect.ImmutableList; import org.apache.cassandra.tcm.Epoch; +import org.apache.cassandra.tcm.serialization.Version; import org.apache.cassandra.tools.NodeProbe; import org.apache.cassandra.tools.nodetool.layout.CassandraUsage; +import picocli.CommandLine.ArgGroup; import picocli.CommandLine.Command; import picocli.CommandLine.Option; import picocli.CommandLine.Parameters; @@ -53,6 +56,7 @@ import static org.apache.cassandra.tcm.CMSOperations.SERVICE_STATE; CMSAdmin.Snapshot.class, CMSAdmin.Unregister.class, CMSAdmin.AbortInitialization.class, + CMSAdmin.DumpClusterMetadata.class, CMSAdmin.DumpDirectory.class, CMSAdmin.DumpLog.class, CMSAdmin.ResumeDropAccordTable.class }) @@ -164,7 +168,7 @@ public class CMSAdmin extends AbstractCommand if (!rf.contains(":")) { if (args.size() > 1) - throw new IllegalArgumentException("Simple placement can only specify a single replication factor accross all data centers"); + throw new IllegalArgumentException("Simple placement can only specify a single replication factor across all data centers"); int parsedRf; try { @@ -295,4 +299,55 @@ public class CMSAdmin extends AbstractCommand probe.getCMSOperationsProxy().resumeDropAccordTable(tableId); } } + + @Command(name = "dump", description = "Dumps cluster metadata into a file") + public static class DumpClusterMetadata extends AbstractCommand + { + @ArgGroup(exclusive = false, multiplicity = "0..1") + DumpOptions dumpOptions; + + static class DumpOptions + { + @Option(names = { "-e", "--epoch" }, paramLabel = "epoch", + description = "Epoch at which cluster metadata should be dumped", required = true) + Long epoch; + + @Option(names = { "-te", "--transform-epoch" }, paramLabel = "transform_epoch", + description = "Force metadata to given X epoch while dumping", required = true) + Long transformEpoch; + + @Option(names = { "-sv", "--serialization-version" }, paramLabel = "serialization_version", + description = "Serialization Version", required = true) + Version version; + } + + @Override + protected void execute(NodeProbe probe) + { + try + { + if (dumpOptions == null) + { + String fileLocation = probe.getCMSOperationsProxy().dumpClusterMetadata(); + printCMSDumpLocation(probe, fileLocation); + } + else + { + String fileLocation = probe.getCMSOperationsProxy().dumpClusterMetadata(dumpOptions.epoch, + dumpOptions.transformEpoch, + dumpOptions.version.name()); + printCMSDumpLocation(probe, fileLocation); + } + } + catch (IOException e) + { + throw new RuntimeException(e); + } + } + + private void printCMSDumpLocation(NodeProbe probe, String fileLocation) + { + probe.output().out.println("Cluster Metadata dump available at " + fileLocation); + } + } } diff --git a/test/distributed/org/apache/cassandra/distributed/test/log/ClusterMetadataDumpTest.java b/test/distributed/org/apache/cassandra/distributed/test/log/ClusterMetadataDumpTest.java index b0f66ff885..7731bde594 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/log/ClusterMetadataDumpTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/log/ClusterMetadataDumpTest.java @@ -25,7 +25,10 @@ import org.junit.Test; import org.apache.cassandra.distributed.Cluster; import org.apache.cassandra.distributed.api.NodeToolResult; import org.apache.cassandra.distributed.test.TestBaseImpl; +import org.apache.cassandra.io.util.File; +import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.tcm.ClusterMetadataService; +import org.apache.cassandra.tcm.membership.NodeVersion; import org.apache.cassandra.tcm.transformations.CustomTransformation; import static org.junit.Assert.assertEquals; @@ -113,4 +116,59 @@ public class ClusterMetadataDumpTest extends TestBaseImpl assertEquals(3, tokensFound); } } + + @Test + public void dumpClusterMetadataTest() throws IOException + { + try (Cluster cluster = init(builder().withNodes(3) + .start())) + { + NodeToolResult res = cluster.get(1).nodetoolResult("cms", "dump"); + res.asserts().success(); + + String stdout = res.getStdout(); + String expectedMsgPrefix = "Cluster Metadata dump available at "; + assertTrue(stdout.contains(expectedMsgPrefix)); + int index = stdout.indexOf(expectedMsgPrefix); + String dumpFile = stdout.substring(index + expectedMsgPrefix.length()).trim(); + assertTrue(new File(dumpFile).exists()); + } + } + + @Test + public void dumpClusterMetadataWithParamsTest() throws IOException + { + try (Cluster cluster = init(builder().withNodes(3) + .start())) + { + long currentEpoch = cluster.get(1).callOnInstance(() -> ClusterMetadata.current().epoch.getEpoch()); + String serVersion = NodeVersion.CURRENT.serializationVersion().toString(); + + NodeToolResult res = cluster.get(1).nodetoolResult("cms", "dump", + "--epoch", "1", + "--transform-epoch", String.valueOf(currentEpoch), + "--serialization-version", serVersion); + res.asserts().success(); + + String stdout = res.getStdout(); + String expectedMsgPrefix = "Cluster Metadata dump available at "; + assertTrue(stdout.contains(expectedMsgPrefix)); + int index = stdout.indexOf(expectedMsgPrefix); + String dumpFile = stdout.substring(index + expectedMsgPrefix.length()).trim(); + assertTrue(new File(dumpFile).exists()); + } + } + + @Test + public void dumpClusterMetadataWithPartialParamsFailsTest() throws IOException + { + try (Cluster cluster = init(builder().withNodes(3) + .start())) + { + // Providing only --epoch without --transform-epoch and --serialization-version + // should be rejected by picocli since all three are required together via @ArgGroup + NodeToolResult res = cluster.get(1).nodetoolResult("cms", "dump", "--epoch", "1"); + res.asserts().failure(); + } + } } diff --git a/test/resources/nodetool/help/cms b/test/resources/nodetool/help/cms index 7cfc82c1c0..eb89133a5c 100644 --- a/test/resources/nodetool/help/cms +++ b/test/resources/nodetool/help/cms @@ -41,6 +41,15 @@ SYNOPSIS [(-u | --username )] cms abortinitialization [--initiator ] + nodetool [(-h | --host )] [(-p | --port )] + [(-pw | --password )] + [(-pwf | --password-file )] + [(-u | --username )] cms dump + [(-e | --epoch )] + [(-sv | --serialization-version + )] + [(-te | --transform-epoch )] + nodetool [(-h | --host )] [(-p | --port )] [(-pw | --password )] [(-pwf | --password-file )] @@ -102,6 +111,14 @@ COMMANDS With --initiator option, The address of the node where `cms initialize` was run. + dump + Dumps cluster metadata into a file + + With --epoch option, Epoch at which cluster metadata should be dumped + + With --transform-epoch option, Force metadata to given X epoch while dumping + + With --serialization-version option, Serialization Version dumpdirectory Dump the directory from the current ClusterMetadata diff --git a/test/resources/nodetool/help/cms$dump b/test/resources/nodetool/help/cms$dump new file mode 100644 index 0000000000..21d0c89f53 --- /dev/null +++ b/test/resources/nodetool/help/cms$dump @@ -0,0 +1,38 @@ +NAME + nodetool cms dump - Dumps cluster metadata into a file + +SYNOPSIS + nodetool [(-h | --host )] [(-p | --port )] + [(-pw | --password )] + [(-pwf | --password-file )] + [(-u | --username )] cms dump + [(-e | --epoch )] + [(-sv | --serialization-version + )] + [(-te | --transform-epoch )] + +OPTIONS + -e , --epoch + Epoch at which cluster metadata should be dumped + + -h , --host + Node hostname or ip address + + -p , --port + Remote jmx agent port number + + -pw , --password + Remote jmx agent password + + -pwf , --password-file + Path to the JMX password file + + -sv , --serialization-version + + Serialization Version + + -te , --transform-epoch + Force metadata to given X epoch while dumping + + -u , --username + Remote jmx agent username diff --git a/test/unit/org/apache/cassandra/tools/CMSOfflineToolTest.java b/test/unit/org/apache/cassandra/tools/CMSOfflineToolTest.java new file mode 100644 index 0000000000..57d87d22e7 --- /dev/null +++ b/test/unit/org/apache/cassandra/tools/CMSOfflineToolTest.java @@ -0,0 +1,1967 @@ +/* + * 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.tools; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Paths; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Random; +import java.util.Set; + +import com.google.common.collect.ImmutableList; +import com.google.common.collect.ImmutableMap; + +import org.junit.Before; +import org.junit.Rule; +import org.junit.Test; +import org.junit.rules.TemporaryFolder; + +import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.dht.IPartitioner; +import org.apache.cassandra.dht.Murmur3Partitioner; +import org.apache.cassandra.dht.Range; +import org.apache.cassandra.dht.Token; +import org.apache.cassandra.io.util.FileInputStreamPlus; +import org.apache.cassandra.io.util.FileOutputStreamPlus; +import org.apache.cassandra.locator.InetAddressAndPort; +import org.apache.cassandra.locator.Replica; +import org.apache.cassandra.schema.DistributedSchema; +import org.apache.cassandra.schema.KeyspaceMetadata; +import org.apache.cassandra.schema.KeyspaceParams; +import org.apache.cassandra.schema.Keyspaces; +import org.apache.cassandra.schema.ReplicationParams; +import org.apache.cassandra.schema.SchemaConstants; +import org.apache.cassandra.service.accord.topology.AccordFastPath; +import org.apache.cassandra.service.accord.topology.AccordStaleReplicas; +import org.apache.cassandra.service.consensus.migration.ConsensusMigrationState; +import org.apache.cassandra.tcm.ClusterMetadata; +import org.apache.cassandra.tcm.ClusterMetadataService; +import org.apache.cassandra.tcm.Epoch; +import org.apache.cassandra.tcm.membership.Directory; +import org.apache.cassandra.tcm.membership.Location; +import org.apache.cassandra.tcm.membership.NodeAddresses; +import org.apache.cassandra.tcm.membership.NodeId; +import org.apache.cassandra.tcm.membership.NodeState; +import org.apache.cassandra.tcm.membership.NodeVersion; +import org.apache.cassandra.tcm.ownership.DataPlacement; +import org.apache.cassandra.tcm.ownership.DataPlacements; +import org.apache.cassandra.tcm.ownership.TokenMap; +import org.apache.cassandra.tcm.ownership.UniformRangePlacement; +import org.apache.cassandra.tcm.sequences.BootstrapAndJoin; +import org.apache.cassandra.tcm.sequences.BootstrapAndReplace; +import org.apache.cassandra.tcm.sequences.InProgressSequences; +import org.apache.cassandra.tcm.sequences.LeaveStreams; +import org.apache.cassandra.tcm.sequences.LockedRanges; +import org.apache.cassandra.tcm.sequences.Move; +import org.apache.cassandra.tcm.sequences.UnbootstrapAndLeave; +import org.apache.cassandra.tcm.serialization.VerboseMetadataSerializer; +import org.apache.cassandra.tcm.serialization.Version; +import org.apache.cassandra.tcm.transformations.PrepareJoin; +import org.apache.cassandra.tcm.transformations.PrepareLeave; +import org.apache.cassandra.tcm.transformations.PrepareMove; +import org.apache.cassandra.tcm.transformations.PrepareReplace; +import org.apache.cassandra.tcm.transformations.Register; + +import static org.apache.cassandra.config.DatabaseDescriptor.getStoragePort; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.prepareJoin; +import static org.assertj.core.api.Assertions.assertThat; + + +public class CMSOfflineToolTest extends OfflineToolUtils +{ + + public static final String DC = "datacenter1"; + @Rule + public final TemporaryFolder temporaryFolder = new TemporaryFolder(); + + private static ClusterMetadata getClusterMetadata(Keyspaces keyspaces, IPartitioner partitioner, Directory directory) + { + DistributedSchema distributedSchema = new DistributedSchema(keyspaces); + + return new ClusterMetadata(Epoch.EMPTY, + partitioner, + distributedSchema, + directory, + new TokenMap(partitioner), + DataPlacements.empty(), + AccordFastPath.EMPTY, + LockedRanges.EMPTY, + InProgressSequences.EMPTY, + ConsensusMigrationState.EMPTY, + ImmutableMap.of(), + AccordStaleReplicas.EMPTY); + } + + @Before + public void setup() + { + DatabaseDescriptor.toolInitialization(); + ClusterMetadataService.initializeForTools(true); + } + + @Test + public void testDefaultCmd() + { + ToolRunner.ToolResult tool = ToolRunner.invokeClass(CMSOfflineTool.class); + + tool.assertOnCleanExit(); + assertThat(tool.getExitCode()).isZero(); + assertThat(tool.getStderr()).isEmpty(); + assertThat(tool.getStdout()).withFailMessage(tool.getStderr()).contains("Usage"); + assertCorrectEnvPostTest(); + } + + @Test + public void testRunCommandThatDoesNotExist() + { + ToolRunner.ToolResult tool = ToolRunner.invokeClass(CMSOfflineTool.class, "cmddoesnotexist"); + assertThat(tool.getExitCode()).isEqualTo(2); + assertCorrectEnvPostTest(); + } + + @Test + public void testRunCommandWithMetadataFileThatDoesnotExist() + { + String metadataFile = temporaryFolder.getRoot().getAbsolutePath() + "/file-does-not-exists.dump"; + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "printdirectory", + "-f", + metadataFile); + + assertThat(result.getExitCode()).isEqualTo(2); + assertCorrectEnvPostTest(); + } + + @Test + public void testAbortbootstrapJoiningNode() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + int nodeToMove = 4; + metadata = startJoining(nodeToMove, metadata); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + String nodeId = String.valueOf(nodeToMove); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "abortbootstrap", + "-f", + metadataFile, + "-id", + nodeId, + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + + ClusterMetadata outMetadata = deserializeMetadata(outputFile); + assertThat(outMetadata.directory.peerIds()).doesNotContain(new NodeId(nodeToMove)); + + assertCorrectEnvPostTest(); + } + + @Test + public void testAbortbootstrapReplacingNode() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + NodeId nodeToReplaceId = new NodeId(3); + NodeId newNodeId = new NodeId(4); + metadata = startReplacing(nodeToReplaceId, newNodeId, metadata); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + String newNodeAddr = metadata.directory.getNodeAddresses(newNodeId).nativeAddress.getHostAddressAndPort(); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "abortbootstrap", + "-f", + metadataFile, + "-ip", + newNodeAddr, + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + + ClusterMetadata outMetadata = deserializeMetadata(outputFile); + assertThat(outMetadata.directory.peerIds()).doesNotContain(newNodeId); + + assertCorrectEnvPostTest(); + } + + @Test + public void testAbortbootstrapNodeGettingReplaced() throws IOException + { + // Assuming that operator unintentionally invokes abortbootstrap on node getting replaced + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + NodeId nodeToReplaceId = new NodeId(3); + NodeId newNodeId = new NodeId(4); + metadata = startReplacing(nodeToReplaceId, newNodeId, metadata); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + String nodeGettingReplaced = metadata.directory.getNodeAddresses(nodeToReplaceId).nativeAddress.getHostAddressAndPort(); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "abortbootstrap", + "-f", + metadataFile, + "-ip", + nodeGettingReplaced, + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("No transformation sequence is in progress for " + nodeGettingReplaced); + + assertCorrectEnvPostTest(); + } + + @Test + public void testAbortbootstrapWhenLeavingSequenceIsInProgress() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(); + NodeId nodeId = new NodeId(4); + metadata = startLeaving(nodeId, metadata); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "abortbootstrap", + "-f", + metadataFile, + "-id", + String.valueOf(nodeId.id()), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("Sequence of kind LEAVE is in progress for node") + .contains("Cannot proceed with this operation"); + + assertCorrectEnvPostTest(); + } + + @Test + public void testAbortbootstrapJoinedNode() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + int nodeId = 4; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "abortbootstrap", + "-f", + metadataFile, + "-id", + String.valueOf(nodeId), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + + assertCorrectEnvPostTest(); + } + + @Test + public void testAbortMoveNoInProgressSequence() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + int nodeId = 4; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "abortmove", + "-f", + metadataFile, + "-id", + String.valueOf(nodeId), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("No transformation sequence is in progress for"); + + assertCorrectEnvPostTest(); + } + + @Test + public void testAbortMoveNonMovingSequence() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(); + NodeId nodeId = new NodeId(4); + metadata = startLeaving(nodeId, metadata); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "abortmove", + "-f", + metadataFile, + "-id", + String.valueOf(nodeId.id()), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("Sequence of kind") + .contains("LEAVE") + .contains("Cannot proceed with this operation"); + + assertCorrectEnvPostTest(); + } + + @Test + public void testAbortMoveMovingNode() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(); + NodeId nodeId = new NodeId(4); + metadata = startMoving(nodeId, metadata); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "abortmove", + "-f", + metadataFile, + "-id", + String.valueOf(nodeId.id()), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + + ClusterMetadata outMetadata = deserializeMetadata(outputFile); + assertThat(outMetadata.inProgressSequences.contains(nodeId)).isFalse(); + + assertCorrectEnvPostTest(); + } + + @Test + public void testAbortDecommissionNoInProgressSequence() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + int nodeId = 4; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "abortdecommission", + "-f", + metadataFile, + "-id", + String.valueOf(nodeId), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("No transformation sequence is in progress for"); + + assertCorrectEnvPostTest(); + } + + @Test + public void testAbortDecommissionNonLeavingSequence() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(); + NodeId nodeId = new NodeId(4); + metadata = startMoving(nodeId, metadata); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "abortdecommission", + "-f", + metadataFile, + "-id", + String.valueOf(nodeId.id()), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("Sequence of kind") + .contains("MOVE") + .contains("Cannot proceed with this operation"); + + assertCorrectEnvPostTest(); + } + + @Test + public void testAbortDecommissionLeavingNode() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(); + NodeId nodeId = new NodeId(4); + metadata = startLeaving(nodeId, metadata); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "abortdecommission", + "-f", + metadataFile, + "-id", + String.valueOf(nodeId.id()), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + + ClusterMetadata outMetadata = deserializeMetadata(outputFile); + assertThat(outMetadata.inProgressSequences.contains(nodeId)).isFalse(); + + assertCorrectEnvPostTest(); + } + + @Test + public void testAssassinateCMSMember() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + // Verify that we have CMS members + assertThat(metadata.fullCMSMembers()).isNotEmpty(); + + // Pick the first CMS member to assassinate + NodeId cmsNodeId = metadata.directory.peerIds().stream().findFirst().orElseThrow(); + assertThat(metadata.isCMSMember(metadata.directory.endpoint(cmsNodeId))).isTrue(); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "assassinate", + "-f", + metadataFile, + "-id", + String.valueOf(cmsNodeId.id()), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + + ClusterMetadata outMetadata = deserializeMetadata(outputFile); + assertThat(outMetadata.directory.peerState(cmsNodeId)).isEqualTo(NodeState.LEFT); + assertThat(outMetadata.fullCMSMemberIds()).doesNotContain(cmsNodeId); + assertThat(outMetadata.fullCMSMembers().size()).isEqualTo(metadata.fullCMSMembers().size()); + + assertCorrectEnvPostTest(); + } + + @Test + public void testAssassinateNodeNonCMSMember() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + // Verify that we have CMS members + assertThat(metadata.fullCMSMembers()).isNotEmpty(); + + // Select non-CMS node + NodeId nonCMSNodeId = metadata.directory.peerIds().stream() + .filter(id -> !metadata.isCMSMember(metadata.directory.endpoint(id))) + .findFirst().orElseThrow(); + assertThat(metadata.isCMSMember(metadata.directory.endpoint(nonCMSNodeId))).isFalse(); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "assassinate", + "-f", + metadataFile, + "-id", + String.valueOf(nonCMSNodeId.id()), + "-o", + outputFile); + + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + ClusterMetadata outMetadata = deserializeMetadata(outputFile); + assertThat(outMetadata.directory.peerState(nonCMSNodeId)).isEqualTo(NodeState.LEFT); + assertCorrectEnvPostTest(); + } + + @Test + public void testAssassinateNodeInvalidNodeId() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + String nodeId = String.valueOf(-1); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "assassinate", + "-f", + metadataFile, + "-id", + nodeId, + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("No node present with id " + nodeId + + " in the given cluster metadata"); + + assertCorrectEnvPostTest(); + } + + @Test + public void testAssassinateMovingNode() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(); + int nodeToMove = 4; + metadata = startMoving(new NodeId(nodeToMove), metadata); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + String nodeId = String.valueOf(nodeToMove); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "assassinate", + "-f", + metadataFile, + "-id", + nodeId, + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + assertThat(result.getStdout()).contains("Cancelling in-progress sequence"); + + ClusterMetadata outMetadata = deserializeMetadata(outputFile); + assertThat(outMetadata.directory.peerState(new NodeId(nodeToMove))).isEqualTo(NodeState.LEFT); + assertCorrectEnvPostTest(); + } + + @Test + public void testAssassinateJoiningNode() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + int nodeToMove = 4; + metadata = startJoining(nodeToMove, metadata); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + String nodeId = String.valueOf(nodeToMove); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "assassinate", + "-f", + metadataFile, + "-id", + nodeId, + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + + // It should suggest to use abortbootstrap for node in joining state + assertThat(result.getStderr()).contains("abortbootstrap"); + + assertCorrectEnvPostTest(); + } + + @Test + public void testAssassinateNodeGettingReplaced() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(); + NodeId nodeToReplaceId = new NodeId(4); + NodeId newNodeId = new NodeId(5); + metadata = startReplacing(nodeToReplaceId, newNodeId, metadata); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + String nodeGettingReplaced = metadata.directory.getNodeAddresses(nodeToReplaceId) + .nativeAddress.getHostAddressAndPort(); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "assassinate", + "-f", + metadataFile, + "-ip", + nodeGettingReplaced, + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("INVALID: Rejecting this plan as it interacts with a range locked"); + + assertCorrectEnvPostTest(); + } + + @Test + public void testAssassinateLeavingNode() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(); + NodeId nodeId = new NodeId(4); + metadata = startLeaving(nodeId, metadata); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + String ip = metadata.directory.endpoint(nodeId).getHostAddressAndPort(); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "assassinate", + "-f", + metadataFile, + "-ip", + ip, + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + assertThat(result.getStdout()).contains("Cancelling in-progress sequence"); + + ClusterMetadata outMetadata = deserializeMetadata(outputFile); + assertThat(outMetadata.directory.peerState(nodeId)).isEqualTo(NodeState.LEFT); + + assertCorrectEnvPostTest(); + } + + @Test + public void testDescribe() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "describe", + "-f", + metadataFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + + int storagePort = getStoragePort(); + + String expectedOutput = + "Cluster Metadata Service:\n" + + "Members: /127.0.0.1:" + storagePort + ",/127.0.0.2:" + storagePort + ",/127.0.0.3:" + storagePort + '\n' + + "Needs reconfiguration: false\n" + + "Service State: LOCAL\n" + + "Epoch: 2\n" + + "Replication factor: ReplicationParams{class=org.apache.cassandra.locator.MetaStrategy, datacenter1=3}\n"; + assertThat(result.getStdout()).isEqualTo(expectedOutput); + + assertCorrectEnvPostTest(); + } + + @Test + public void testDescribeRawString() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "print", + "-f", + metadataFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(result.getStdout()).isNotEmpty().startsWith("ClusterMetadata{"); + + assertCorrectEnvPostTest(); + } + + @Test + public void testResetCMSUsingNodeId() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out" + new Random().nextLong() + ".dump"; + NodeId nodeId = new NodeId(1); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "resetcms", + "-f", + metadataFile, + "-id", + String.valueOf(nodeId.id()), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(result.getStderr()).isEmpty(); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + + ClusterMetadata outMetadata = deserializeMetadata(outputFile); + InetAddressAndPort nodeAddress = metadata.directory.getNodeAddresses(nodeId).broadcastAddress; + assertThat(outMetadata.isCMSMember(nodeAddress)).isTrue(); + assertThat(outMetadata.fullCMSMembers().size()).isEqualTo(1); + + assertCorrectEnvPostTest(); + } + + @Test + public void testResetCMS() throws IOException + { + assertResetCMS("127.0.0.1:" + getStoragePort()); + } + + @Test + public void testResetCMSUsingIp() throws IOException + { + assertResetCMS("127.0.0.1"); + } + + private void assertResetCMS(String ipAddress) throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out" + new Random().nextLong() + ".dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "resetcms", + "-f", + metadataFile, + "-ip", + ipAddress, + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(result.getStderr()).isEmpty(); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + + ClusterMetadata outMetadata = deserializeMetadata(outputFile); + assertThat(outMetadata).isNotNull(); + + // Check given ip address is added to CMS members + InetAddressAndPort candidate = InetAddressAndPort.getByNameUnchecked(ipAddress); + assertThat(outMetadata.isCMSMember(candidate)).isTrue(); + assertThat(outMetadata.fullCMSMembers().size()).isEqualTo(1); + + assertCorrectEnvPostTest(); + } + + @Test + public void testResetCMSUsingIpAddressThatDoesNotExist() throws IOException + { + // Node that doesn't exist in Cluster Metadata can be added as CMS member + String ipAddress = "127.0.0.55:" + getStoragePort(); + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "resetcms", + "-f", + metadataFile, + "-ip", + ipAddress, + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isNotEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("No node present").contains(ipAddress); + + assertCorrectEnvPostTest(); + } + + @Test + public void testResetCMSInvalidIpAddress() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + String invalidIpAddress = "/127.0.0."; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "resetcms", + "-f", + metadataFile, + "-ip", + invalidIpAddress, + "-o", + outputFile); + + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(2); + assertThat(result.getStderr()).contains("java.net.UnknownHostException"); + + assertCorrectEnvPostTest(); + } + + @Test + public void testResetCMSInvalidSerializationVersion() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + String invalidIpAddress = "127.0.0.3"; + String serializationVersion = "-1"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "resetcms", + "-f", + metadataFile, + "-sv", + serializationVersion, + "-ip", + invalidIpAddress, + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("Invalid value for option '--serialization-version'"); + assertCorrectEnvPostTest(); + } + + @Test + public void testMoveToken() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(16); + String nodeToMove = "127.0.0.4:" + getStoragePort(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "move", + "-f", + metadataFile, + "-ip", + nodeToMove, + "-t", + String.valueOf(metadata.partitioner.getRandomToken()), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("This node has more than one token and cannot be moved thusly."); + + assertCorrectEnvPostTest(); + } + + @Test + public void testMoveTokenVnodes() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(4); + metadata = addNewNode(metadata, 4, Set.of(metadata.partitioner.getRandomToken())); + String nodeToMove = "127.0.0.4:" + getStoragePort(); + InetAddressAndPort newNodeInetAddress = InetAddressAndPort.getByNameUnchecked(nodeToMove); + NodeId nodeId = metadata.directory.peerId(InetAddressAndPort.getByName(nodeToMove)); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + Token desiredToken = metadata.partitioner.getRandomToken(); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "move", + "-f", + metadataFile, + "-ip", + nodeToMove, + "-t", + String.valueOf(desiredToken.getTokenValue()), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + + ClusterMetadata outMetadata = deserializeMetadata(outputFile); + assertThat(outMetadata).isNotNull(); + assertThat(outMetadata.directory.peerId(newNodeInetAddress)).isNotNull(); + assertThat(outMetadata.tokenMap.tokens(nodeId)).contains(desiredToken); + + assertCorrectEnvPostTest(); + } + + @Test + public void testMoveTokenWithoutToken() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(); + String nodeToMove = "127.0.0.4:" + getStoragePort(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "move", + "-f", + metadataFile, + "-ip", + nodeToMove, + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("Token required").contains("MOVE"); + assertCorrectEnvPostTest(); + } + + @Test + public void testMoveTokenForNonExistingNode() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + String unknownNodeIpWithPort = "127.0.0.5:" + getStoragePort(); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + Token tokenToAssign = metadata.partitioner.getRandomToken(); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "move", + "-f", + metadataFile, + "-ip", + unknownNodeIpWithPort, + "-t", + String.valueOf(tokenToAssign.getTokenValue()), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(2); + assertThat(result.getStderr()).contains("No node present with ip address") + .contains(unknownNodeIpWithPort); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + + assertCorrectEnvPostTest(); + } + + @Test + public void testMoveInvalidIpAddress() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + String invalidIp = "127.0.0."; + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + Token tokenToAssign = metadata.partitioner.getRandomToken(); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "move", + "-f", + metadataFile, + "-ip", + invalidIp, + "-t", + String.valueOf(tokenToAssign.getTokenValue()), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(2); + assertThat(result.getStderr()).contains("java.net.UnknownHostException"); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + + assertCorrectEnvPostTest(); + } + + @Test + public void testMoveNodeToInvalidToken() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + String newNodeIpWithPort = "127.0.0.1:" + getStoragePort(); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + String invalidToken = "somegibberishinvalidtoken"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "move", + "-f", + metadataFile, + "-ip", + newNodeIpWithPort, + "-t", + invalidToken, + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + + assertCorrectEnvPostTest(); + } + + @Test + public void testMoveNodeToSomeOtherNodeToken() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + String newNodeIpWithPort = "127.0.0.1:" + getStoragePort(); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + NodeId secondNode = new NodeId(2); + ImmutableList secondNodeTokenList = metadata.tokenMap.tokens(secondNode); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "move", + "-f", + metadataFile, + "-ip", + newNodeIpWithPort, + "-t", + secondNodeTokenList.get(0).toString(), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("is already owned by node " + secondNode.id()); + + assertCorrectEnvPostTest(); + } + + @Test + public void testMoveNodeWhenMoveInProgressToAnotherToken() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + NodeId nodeId = new NodeId(2); + Token randomToken = metadata.partitioner.getRandomToken(); + Set toTokenSet = Set.of(randomToken); + metadata = startMoving(nodeId, metadata, toTokenSet); + String metadataFile = dumpMetadata(metadata); + String newNodeIpWithPort = "127.0.0." + nodeId.id() + ':' + getStoragePort(); + + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "move", + "-f", + metadataFile, + "-ip", + newNodeIpWithPort, + "-t", + randomToken.toString(), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + + ClusterMetadata outMetadata = deserializeMetadata(outputFile); + assertThat(outMetadata.tokenMap.tokens(nodeId)).isEqualTo(ImmutableList.of(randomToken)); + + assertCorrectEnvPostTest(); + } + + @Test + public void testMoveNodeFinishInProgress() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + NodeId nodeId = new NodeId(2); + metadata = startMoving(nodeId, metadata); + String metadataFile = dumpMetadata(metadata); + String newNodeIpWithPort = "127.0.0." + nodeId.id() + ':' + getStoragePort(); + + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "move", + "-f", + metadataFile, + "-ip", + newNodeIpWithPort, + "-t", + metadata.partitioner.getRandomToken().toString(), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("Move in progress for another token(s)"); + + assertCorrectEnvPostTest(); + } + + @Test + public void testForceJoin() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String newNodeIpWithPort = "127.0.0.5:" + getStoragePort(); + InetAddressAndPort newNodeInetAddress = InetAddressAndPort.getByNameUnchecked(newNodeIpWithPort); + Directory newDirectory = metadata.directory.with(new NodeAddresses(newNodeInetAddress), + new Location("datacenter1", "rack4"), + NodeVersion.CURRENT); + metadata = updateMetadata(metadata, newDirectory); + + NodeId newNodeId = metadata.directory.peerId(newNodeInetAddress); + assertThat(metadata.directory.states.get(newNodeId)).isEqualTo(NodeState.REGISTERED); + + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "forcejoin", + "-f", + metadataFile, + "-id", + String.valueOf(newNodeId.id()), + "-t", + metadata.partitioner.getRandomToken().toString(), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + + ClusterMetadata outMetadata = deserializeMetadata(outputFile); + assertThat(outMetadata).isNotNull(); + assertThat(outMetadata.directory.peerId(newNodeInetAddress)).isNotNull(); + assertThat(outMetadata.directory.states.get(newNodeId)).isEqualTo(NodeState.JOINED); + assertThat(outMetadata.tokenMap.tokens(newNodeId)).isNotEmpty(); + + assertCorrectEnvPostTest(); + } + + @Test + public void testForceJoinNewNodeWithoutTokens() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String newNodeIpWithPort = "127.0.0.5:" + getStoragePort(); + InetAddressAndPort newNodeInetAddress = InetAddressAndPort.getByNameUnchecked(newNodeIpWithPort); + Directory newDirectory = metadata.directory.with(new NodeAddresses(newNodeInetAddress), + new Location("datacenter1", "rack4"), + NodeVersion.CURRENT); + metadata = updateMetadata(metadata, newDirectory); + + NodeId newNodeId = metadata.directory.peerId(newNodeInetAddress); + assertThat(metadata.directory.states.get(newNodeId)).isEqualTo(NodeState.REGISTERED); + + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "forcejoin", + "-f", + metadataFile, + "-id", + String.valueOf(newNodeId.id()), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("Tokens must be provided to force join a node."); + + assertCorrectEnvPostTest(); + } + + @Test + public void testForceJoinUnknownNode() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String newNodeIpWithPort = "127.0.0.5:" + getStoragePort(); + // New node is not registered, it should fail + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "forcejoin", + "-f", + metadataFile, + "-ip", + newNodeIpWithPort, + "-t", + metadata.partitioner.getRandomToken().toString(), + "-t", + metadata.partitioner.getRandomToken().toString(), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isNotEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains(newNodeIpWithPort) + .contains("No node present with ip address"); + + assertCorrectEnvPostTest(); + } + + @Test + public void testForceJoinMovingNode() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(); + NodeId nodeId = new NodeId(4); + metadata = startMoving(nodeId, metadata); + assertThat(metadata.directory.peerState(nodeId)).isEqualTo(NodeState.MOVING); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "forcejoin", + "-f", + metadataFile, + "-id", + String.valueOf(nodeId.id()), + "-t", + metadata.partitioner.getRandomToken().toString(), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("Another sequence of kind MOVE is in progress for node " + nodeId.id() + + ". Cannot proceed with force join."); + + assertCorrectEnvPostTest(); + } + + @Test + public void testForceJoinAlreadyJoinedNode() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(); + NodeId nodeId = new NodeId(4); + assertThat(metadata.directory.peerState(nodeId)).isEqualTo(NodeState.JOINED); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "forcejoin", + "-f", + metadataFile, + "-id", + String.valueOf(nodeId.id()), + "-t", + metadata.partitioner.getRandomToken().toString(), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("Node " + nodeId.id() + " is already in JOINED state."); + + assertCorrectEnvPostTest(); + } + + @Test + public void testForceJoinLeftNode() throws IOException + { + ClusterMetadata metadata = getFourNodeMetadata(); + int id = 4; + NodeId nodeId = new NodeId(id); + metadata = startLeaving(nodeId, metadata); + UnbootstrapAndLeave unbootstrapAndLeave = (UnbootstrapAndLeave) metadata.inProgressSequences.get(nodeId); + metadata = unbootstrapAndLeave.applyTo(metadata).success().metadata; + assertThat(metadata.directory.peerState(nodeId)).isEqualTo(NodeState.LEFT); + + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "forcejoin", + "-f", + metadataFile, + "-id", + String.valueOf(nodeId.id()), + "-t", + metadata.partitioner.getRandomToken().toString(), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + + ClusterMetadata outMetadata = deserializeMetadata(outputFile); + assertThat(outMetadata.directory.peerState(nodeId)).isEqualTo(NodeState.JOINED); + assertCorrectEnvPostTest(); + } + + @Test + public void testForceJoinJoiningNodeTokenMismatch() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + int nodeNum = 4; + metadata = startJoining(nodeNum, metadata); + NodeId nodeId = new NodeId(nodeNum); + assertThat(metadata.directory.peerState(nodeId)).isEqualTo(NodeState.BOOTSTRAPPING); + + BootstrapAndJoin bootstrapAndJoin = (BootstrapAndJoin) metadata.inProgressSequences.get(nodeId); + Set sequenceTokens = bootstrapAndJoin.finishJoin.tokens; + + // Pick a token that is different from the in-progress sequence tokens + Token differentToken = metadata.partitioner.getRandomToken(); + while (sequenceTokens.contains(differentToken)) + differentToken = metadata.partitioner.getRandomToken(); + + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "forcejoin", + "-f", + metadataFile, + "-id", + String.valueOf(nodeId.id()), + "-t", + differentToken.toString(), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("do not match with") + .contains("in progress BootstrapAndJoin sequence tokens"); + + assertCorrectEnvPostTest(); + } + + @Test + public void testForceJoinJoiningNodeTokenMatch() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + int nodeNum = 4; + metadata = startJoining(nodeNum, metadata); + NodeId nodeId = new NodeId(nodeNum); + assertThat(metadata.directory.peerState(nodeId)).isEqualTo(NodeState.BOOTSTRAPPING); + + BootstrapAndJoin bootstrapAndJoin = (BootstrapAndJoin) metadata.inProgressSequences.get(nodeId); + Set sequenceTokens = bootstrapAndJoin.finishJoin.tokens; + + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + // Build args: one -t per token in the sequence + List args = new ArrayList<>(List.of("forcejoin", "-f", metadataFile, + "-id", String.valueOf(nodeId.id()), + "-o", outputFile)); + for (Token t : sequenceTokens) + { + args.add("-t"); + args.add(t.toString()); + } + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + args.toArray(new String[0])); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + + ClusterMetadata outMetadata = deserializeMetadata(outputFile); + assertThat(outMetadata.directory.peerState(nodeId)).isEqualTo(NodeState.JOINED); + assertThat(outMetadata.tokenMap.tokens(nodeId)).containsAll(sequenceTokens); + + assertCorrectEnvPostTest(); + } + + @Test + public void testForceJoinNodeThatDoesNotExist() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + String nodeId = "-1"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "forcejoin", + "-f", + metadataFile, + "-id", + nodeId, + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).isNotEmpty(); + + assertCorrectEnvPostTest(); + } + + @Test + public void testForceJoinWithSerializationVersion() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String newNodeIpWithPort = "127.0.0.5:" + getStoragePort(); + InetAddressAndPort newNodeInetAddress = InetAddressAndPort.getByNameUnchecked(newNodeIpWithPort); + Directory newDirectory = metadata.directory.with(new NodeAddresses(newNodeInetAddress), + new Location("datacenter1", "rack4"), + NodeVersion.CURRENT); + metadata = updateMetadata(metadata, newDirectory); + + NodeId newNodeId = metadata.directory.peerId(newNodeInetAddress); + assertThat(metadata.directory.states.get(newNodeId)).isEqualTo(NodeState.REGISTERED); + + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + // Use V7 — a known valid version older than the current (V8) + Version targetVersion = Version.V7; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "forcejoin", + "-f", + metadataFile, + "-sv", + targetVersion.toString(), + "-id", + String.valueOf(newNodeId.id()), + "-t", + metadata.partitioner.getRandomToken().toString(), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + + // Verify that the output file was serialized using the requested version + try (FileInputStreamPlus fisp = new FileInputStreamPlus(outputFile)) + { + int versionInt = fisp.readUnsignedVInt32(); + assertThat(Version.fromInt(versionInt)).isEqualTo(targetVersion); + } + + // Verify the node was force-joined in the output metadata + ClusterMetadata outMetadata = deserializeMetadata(outputFile); + assertThat(outMetadata.directory.peerState(newNodeId)).isEqualTo(NodeState.JOINED); + + assertCorrectEnvPostTest(); + } + + @Test + public void testForceJoinInvalidSerializationVersion() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String newNodeIpWithPort = "127.0.0.5:" + getStoragePort(); + InetAddressAndPort newNodeInetAddress = InetAddressAndPort.getByNameUnchecked(newNodeIpWithPort); + Directory newDirectory = metadata.directory.with(new NodeAddresses(newNodeInetAddress), + new Location("datacenter1", "rack4"), + NodeVersion.CURRENT); + metadata = updateMetadata(metadata, newDirectory); + + NodeId newNodeId = metadata.directory.peerId(newNodeInetAddress); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + String invalidSerializationVersion = "-1"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "forcejoin", + "-f", + metadataFile, + "-sv", + invalidSerializationVersion, + "-id", + String.valueOf(newNodeId.id()), + "-t", + metadata.partitioner.getRandomToken().toString(), + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("Invalid value for option '--serialization-version'"); + + assertCorrectEnvPostTest(); + } + + @Test + public void testPrintDataPlacements() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String keyspaceName = "ks1"; + KeyspaceParams ksParams = KeyspaceParams.create(true, + Map.of("class", "NetworkTopologyStrategy", + "datacenter1", "3")); + Keyspaces keyspaces = Keyspaces.none().with(KeyspaceMetadata.create(keyspaceName, ksParams)); + DistributedSchema newSchema = new DistributedSchema(keyspaces); + metadata = updateMetadata(metadata, newSchema); + + InetAddressAndPort inetAddressAndPort = InetAddressAndPort.getByNameUnchecked("127.0.0.1:" + getStoragePort()); + InetAddressAndPort i2 = InetAddressAndPort.getByNameUnchecked("127.0.0.2:" + getStoragePort()); + InetAddressAndPort i3 = InetAddressAndPort.getByNameUnchecked("127.0.0.3:" + getStoragePort()); + + Range tokenRange = new Range<>(metadata.partitioner.getMinimumToken(), metadata.partitioner.getTokenFactory().fromString("0")); + + DataPlacement dataPlacement = DataPlacement.builder() + .withWriteReplica(Epoch.FIRST, Replica.fullReplica(inetAddressAndPort, tokenRange)) + .withWriteReplica(Epoch.FIRST, Replica.fullReplica(i2, tokenRange)) + .withWriteReplica(Epoch.FIRST, Replica.fullReplica(i3, tokenRange)) + .withReadReplica(Epoch.FIRST, Replica.fullReplica(inetAddressAndPort, tokenRange)) + .withReadReplica(Epoch.FIRST, Replica.fullReplica(i2, tokenRange)) + .withReadReplica(Epoch.FIRST, Replica.fullReplica(i3, tokenRange)) + .build(); + + DataPlacements dataPlacements = DataPlacements.builder(1) + .with(ReplicationParams.fromMap(Map.of("class", "NetworkTopologyStrategy", + "datacenter1", "3")), + dataPlacement) + .build(); + metadata = updateMetadata(metadata, dataPlacements); + String metadataFile = dumpMetadata(metadata); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "printdataplacements", + "-f", + metadataFile, + "-ks", + keyspaceName); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(result.getStdout()).isNotEmpty(); + assertThat(result.getStderr()).isNullOrEmpty(); + + String stdout = result.getStdout(); + assertThat(stdout).isEqualTo( + "Token Range Type Endpoints\n" + + "(-9223372036854775808,0] read /127.0.0.1:" + getStoragePort() + ", /127.0.0.2:" + getStoragePort() + ", /127.0.0.3:" + getStoragePort() + '\n' + + "(-9223372036854775808,0] write /127.0.0.1:" + getStoragePort() + ", /127.0.0.2:" + getStoragePort() + ", /127.0.0.3:" + getStoragePort() + '\n'); + + assertCorrectEnvPostTest(); + } + + @Test + public void testPrintDataPlacementsKeyspaceDoesNotExist() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String keyspaceName = "ks1"; + String metadataFile = dumpMetadata(metadata); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "printdataplacements", + "-f", + metadataFile, + "-ks", + keyspaceName); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(2); + assertThat(result.getStderr()).isNotEmpty(); + assertThat(result.getStderr()).contains("Keyspace " + keyspaceName + " not found in cluster metadata"); + + assertCorrectEnvPostTest(); + } + + @Test + public void testPrintDirectory() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "printdirectory", + "-f", + metadataFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(result.getStdout()).isNotEmpty(); + assertThat(result.getStderr()).isNullOrEmpty(); + + String stdout = result.getStdout(); + + String expectedSerializationVersion = NodeVersion.CURRENT.serializationVersion().toString(); + assertThat(stdout).isEqualTo( + "NodeId: 1\n" + + " rack rack1\n" + + " local_port " + getStoragePort() + '\n' + + " broadcast_port " + getStoragePort() + '\n' + + " host_id 6d194555-f6eb-41d0-c000-000000000001\n" + + " broadcast_address /127.0.0.1\n" + + " native_address /127.0.0.1\n" + + " native_port " + getStoragePort() + '\n' + + " local_address /127.0.0.1\n" + + " state JOINED\n" + + " serialization_version " + expectedSerializationVersion + '\n' + + " cassandra_version " + metadata.directory.version(new NodeId(1)).cassandraVersion + '\n' + + " dc datacenter1\n" + + " is_cms_member true\n" + + "NodeId: 2\n" + + " rack rack2\n" + + " local_port " + getStoragePort() + '\n' + + " broadcast_port " + getStoragePort() + '\n' + + " host_id 6d194555-f6eb-41d0-c000-000000000002\n" + + " broadcast_address /127.0.0.2\n" + + " native_address /127.0.0.2\n" + + " native_port " + getStoragePort() + '\n' + + " local_address /127.0.0.2\n" + + " state JOINED\n" + + " serialization_version " + expectedSerializationVersion + '\n' + + " cassandra_version " + metadata.directory.version(new NodeId(2)).cassandraVersion + '\n' + + " dc datacenter1\n" + + " is_cms_member true\n" + + "NodeId: 3\n" + + " rack rack3\n" + + " local_port " + getStoragePort() + '\n' + + " broadcast_port " + getStoragePort() + '\n' + + " host_id 6d194555-f6eb-41d0-c000-000000000003\n" + + " broadcast_address /127.0.0.3\n" + + " native_address /127.0.0.3\n" + + " native_port " + getStoragePort() + '\n' + + " local_address /127.0.0.3\n" + + " state JOINED\n" + + " serialization_version " + expectedSerializationVersion + '\n' + + " cassandra_version " + metadata.directory.version(new NodeId(3)).cassandraVersion + '\n' + + " dc datacenter1\n" + + " is_cms_member true\n" + ); + + assertCorrectEnvPostTest(); + } + + private String dumpMetadata(ClusterMetadata metadata) throws IOException + { + String tempFile = temporaryFolder.newFile().getAbsolutePath(); + try (FileOutputStreamPlus out = new FileOutputStreamPlus(tempFile)) + { + VerboseMetadataSerializer.serialize(ClusterMetadata.serializer, + metadata, + out, + NodeVersion.CURRENT.serializationVersion()); + } + return tempFile; + } + + /** + * Creates a three-node cluster metadata for testing. + */ + + private ClusterMetadata getThreeNodeClusterMetadata() + { + return getThreeNodeClusterMetadata(1); + } + + private ClusterMetadata getThreeNodeClusterMetadata(int tokenSize) + { + IPartitioner partitioner = Murmur3Partitioner.instance; + NodeId nodeId1 = new NodeId(1); + NodeId nodeId2 = new NodeId(2); + NodeId nodeId3 = new NodeId(3); + + InetAddressAndPort addr1 = InetAddressAndPort.getByNameUnchecked("127.0.0.1:" + getStoragePort()); + InetAddressAndPort addr2 = InetAddressAndPort.getByNameUnchecked("127.0.0.2:" + getStoragePort()); + InetAddressAndPort addr3 = InetAddressAndPort.getByNameUnchecked("127.0.0.3:" + getStoragePort()); + + NodeVersion nodeVersion = NodeVersion.CURRENT; + + Directory directory = new Directory() + .unsafeWithNodeForTesting(nodeId1, new NodeAddresses(addr1), + new Location(DC, "rack1"), nodeVersion) + .unsafeWithNodeForTesting(nodeId2, new NodeAddresses(addr2), + new Location(DC, "rack2"), nodeVersion) + .unsafeWithNodeForTesting(nodeId3, new NodeAddresses(addr3), + new Location(DC, "rack3"), nodeVersion) + .withRackAndDC(nodeId1) + .withRackAndDC(nodeId2) + .withRackAndDC(nodeId3); + + KeyspaceParams metaKsParams = KeyspaceParams.create(true, + Map.of("class", "MetaStrategy", "datacenter1", "3")); + KeyspaceMetadata metaKeyspace = KeyspaceMetadata.create(SchemaConstants.METADATA_KEYSPACE_NAME, metaKsParams); + KeyspaceMetadata normalKeyspace = KeyspaceMetadata.create("ks", KeyspaceParams.simple(3)); + Keyspaces keyspaces = Keyspaces.none().with(metaKeyspace).with(normalKeyspace); + + ClusterMetadata clusterMetadata = getClusterMetadata(keyspaces, partitioner, directory); + + + ClusterMetadata metadata = clusterMetadata + .transformer() + .with(directory) + .join(nodeId1) + .proposeToken(nodeId1, getRandomTokens(partitioner, tokenSize)) + .join(nodeId2) + .proposeToken(nodeId2, getRandomTokens(partitioner, tokenSize)) + .join(nodeId3) + .proposeToken(nodeId3, getRandomTokens(partitioner, tokenSize)) + .build().metadata; + + // Create replicas for the metadata keyspace on all three nodes + ReplicationParams metaParams = ReplicationParams.ntsMeta(Collections.singletonMap(DC, 3)); + DataPlacements placements = DataPlacements.empty().unbuild() + .with(metaParams, getCMSMemberPlacement(metadata, List.of(addr1, addr2, addr3))) + .with(ReplicationParams.simple(3), getKeyspacePlacement(metadata, normalKeyspace)) + .build(); + + return updateMetadata(metadata, placements); + } + + List getRandomTokens(IPartitioner partitioner, int size) + { + List tokens = new ArrayList<>(size); + for (int i = 0; i < size; i++) + { + tokens.add(partitioner.getRandomToken()); + } + + return tokens; + } + + DataPlacement getCMSMemberPlacement(ClusterMetadata clusterMetadata, List inetAddressAndPorts) + { + IPartitioner partitioner = clusterMetadata.partitioner; + // Create replicas for the metadata keyspace on all three nodes + Range fullRange = new Range<>(partitioner.getMinimumToken(), partitioner.getMinimumToken()); + DataPlacement.Builder placementBuilder = DataPlacement.builder(); + for (InetAddressAndPort addr : inetAddressAndPorts) + { + Replica replica = Replica.fullReplica(addr, fullRange); + placementBuilder.withReadReplica(Epoch.EMPTY, replica) + .withWriteReplica(Epoch.EMPTY, replica); + } + + return placementBuilder.build(); + } + + DataPlacement getKeyspacePlacement(ClusterMetadata metadata, KeyspaceMetadata keyspaceMetadata) + { + UniformRangePlacement uniformRangePlacement = new UniformRangePlacement(); + List> tokenRanges = uniformRangePlacement.calculateRanges(metadata.tokenMap); + + return keyspaceMetadata.replicationStrategy.calculateDataPlacement(metadata.epoch, tokenRanges, metadata); + } + + ClusterMetadata getFourNodeMetadata() + { + return getFourNodeMetadata(1); + } + + ClusterMetadata getFourNodeMetadata(int tokenSize) + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + return addNewNode(metadata, 4, new HashSet<>(getRandomTokens(metadata.partitioner, tokenSize))); + } + + @SuppressWarnings("SameParameterValue") + ClusterMetadata addNewNode(ClusterMetadata prev, int newNodeId, Set tokens) + { + NodeId nodeId = new NodeId(newNodeId); + Directory directory = prev.directory + .unsafeWithNodeForTesting(nodeId, getNodeAddresses(newNodeId), + new Location(DC, "rack" + newNodeId), + NodeVersion.CURRENT) + .withRackAndDC(nodeId); + + ClusterMetadata metadata = prev.transformer().with(directory).join(nodeId).proposeToken(nodeId, tokens).build().metadata; + DataPlacements dataPlacements = new UniformRangePlacement() + .calculatePlacements(metadata.epoch.nextEpoch(), + metadata, + metadata.schema.getKeyspaces()); + return updateMetadata(metadata, dataPlacements); + } + + NodeAddresses getNodeAddresses(int num) + { + String address = "127.0.0." + num; + InetAddressAndPort addressAndPort = InetAddressAndPort.getByNameUnchecked(address + ':' + getStoragePort()); + return new NodeAddresses(addressAndPort); + } + + ClusterMetadata startMoving(NodeId nodeId, ClusterMetadata metadata) + { + return startMoving(nodeId, metadata, Set.of(metadata.partitioner.getRandomToken())); + } + + ClusterMetadata startMoving(NodeId nodeId, ClusterMetadata metadata, Set tokens) + { + PrepareMove prepareMove = new PrepareMove(nodeId, tokens, new UniformRangePlacement(), false); + metadata = prepareMove.execute(metadata).success().metadata; + + Move moveSequence = (Move) metadata.inProgressSequences.get(nodeId); + return moveSequence.startMove.execute(metadata).success().metadata; + } + + ClusterMetadata startJoining(int nodeNum, ClusterMetadata metadata) + { + NodeId nodeId = new NodeId(nodeNum); + + if (!metadata.directory.peerIds().contains(nodeId)) + { + InetAddressAndPort addr = InetAddressAndPort.getByNameUnchecked("127.0.0." + nodeNum + ':' + getStoragePort()); + + Location location = new Location(DC, "rack" + nodeNum); + Register register = new Register(new NodeAddresses(addr), location, NodeVersion.CURRENT); + metadata = register.execute(metadata).success().metadata; + } + + PrepareJoin prepareJoin = prepareJoin(nodeId); + + ClusterMetadata updatedMetadata = prepareJoin.execute(metadata).success().metadata; + BootstrapAndJoin bootstrapAndJoin = (BootstrapAndJoin) updatedMetadata.inProgressSequences.get(nodeId); + return bootstrapAndJoin.startJoin.execute(updatedMetadata).success().metadata; + } + + ClusterMetadata startReplacing(NodeId oldNodeId, NodeId newNodeId, ClusterMetadata clusterMetadata) + { + Register register = new Register(getNodeAddresses(newNodeId.id()), + new Location(DC, "rack" + newNodeId.id()), + NodeVersion.CURRENT); + clusterMetadata = register.execute(clusterMetadata).success().metadata; + PrepareReplace prepareReplace = new PrepareReplace(oldNodeId, newNodeId, new UniformRangePlacement(), + true, false); + ClusterMetadata updatedMetadata = prepareReplace.execute(clusterMetadata).success().metadata; + BootstrapAndReplace replaceSequence = (BootstrapAndReplace) updatedMetadata.inProgressSequences.get(newNodeId); + + updatedMetadata = replaceSequence.startReplace.execute(updatedMetadata).success().metadata; + return updatedMetadata; + } + + ClusterMetadata startLeaving(NodeId nodeId, ClusterMetadata metadata) + { + PrepareLeave prepareLeave = new PrepareLeave(nodeId, true, new UniformRangePlacement(), LeaveStreams.Kind.UNBOOTSTRAP); + ClusterMetadata updatedMetadata = prepareLeave.execute(metadata).success().metadata; + UnbootstrapAndLeave unbootstrapAndLeave = (UnbootstrapAndLeave) updatedMetadata.inProgressSequences.get(nodeId); + return unbootstrapAndLeave.startLeave.execute(updatedMetadata).success().metadata; + } + + private ClusterMetadata updateMetadata(ClusterMetadata metadata, Directory directory) + { + return metadata.transformer().with(directory).build().metadata; + } + + private ClusterMetadata updateMetadata(ClusterMetadata metadata, DataPlacements dataPlacements) + { + return metadata.transformer().with(dataPlacements).build().metadata; + } + + private ClusterMetadata updateMetadata(ClusterMetadata metadata, DistributedSchema newSchema) + { + return metadata.transformer().with(newSchema).build().metadata; + } + + ClusterMetadata deserializeMetadata(String metadataDumpFile) throws IOException + { + return ClusterMetadataService.deserializeClusterMetadata(metadataDumpFile); + } + + // -- getSerializationVersion tests -- + + @Test + public void testDefaultSerializationVersionUsesMetadataVersion() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "resetcms", + "-f", + metadataFile, + "-ip", + "127.0.0.1", + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + + try (FileInputStreamPlus fisp = new FileInputStreamPlus(outputFile)) + { + int versionInt = fisp.readUnsignedVInt32(); + assertThat(Version.fromInt(versionInt)).isEqualTo(metadata.directory.commonSerializationVersion); + } + + assertCorrectEnvPostTest(); + } + + @Test + public void testUserSpecifiedSerializationVersionTakesPrecedence() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + Version userVersion = Version.V7; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "resetcms", + "-f", + metadataFile, + "-sv", + userVersion.toString(), + "-ip", + "127.0.0.1", + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + + try (FileInputStreamPlus fisp = new FileInputStreamPlus(outputFile)) + { + int versionInt = fisp.readUnsignedVInt32(); + assertThat(Version.fromInt(versionInt)).isEqualTo(userVersion); + } + + // V7 is older than metadata version (V8), so a warning should be emitted + assertThat(result.getStderr()).contains("WARNING"); + + assertCorrectEnvPostTest(); + } + + @Test + public void testSerializationVersionNewerThanCurrentFails() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + + // V8 is the current version; passing a version newer than that should fail. + // UNKNOWN has value Integer.MAX_VALUE, which is guaranteed to be newer. + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "resetcms", + "-f", + metadataFile, + "-sv", + Version.UNKNOWN.toString(), + "-ip", + "127.0.0.1", + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStdout()).isEqualTo(2); + assertThat(Files.exists(Paths.get(outputFile))).isFalse(); + assertThat(result.getStderr()).contains("is older than the target serialization version"); + + assertCorrectEnvPostTest(); + } + + @Test + public void testSerializationVersionSameAsMetadataVersion() throws IOException + { + ClusterMetadata metadata = getThreeNodeClusterMetadata(); + String metadataFile = dumpMetadata(metadata); + String outputFile = temporaryFolder.getRoot() + "/metadata-out.dump"; + Version metadataVersion = metadata.directory.commonSerializationVersion; + + ToolRunner.ToolResult result = ToolRunner.invokeClass(CMSOfflineTool.class, + "resetcms", + "-f", + metadataFile, + "-sv", + metadataVersion.toString(), + "-ip", + "127.0.0.1", + "-o", + outputFile); + + assertThat(result.getExitCode()).withFailMessage(result.getStderr()).isEqualTo(0); + assertThat(Files.exists(Paths.get(outputFile))).isTrue(); + + try (FileInputStreamPlus fisp = new FileInputStreamPlus(outputFile)) + { + int versionInt = fisp.readUnsignedVInt32(); + assertThat(Version.fromInt(versionInt)).isEqualTo(metadataVersion); + } + + // Same version as metadata, no warning expected + assertThat(result.getStderr()).doesNotContain("WARNING"); + + assertCorrectEnvPostTest(); + } +} diff --git a/test/unit/org/apache/cassandra/tools/OfflineToolUtils.java b/test/unit/org/apache/cassandra/tools/OfflineToolUtils.java index 9c950fda80..1b486e7f29 100644 --- a/test/unit/org/apache/cassandra/tools/OfflineToolUtils.java +++ b/test/unit/org/apache/cassandra/tools/OfflineToolUtils.java @@ -108,7 +108,9 @@ public abstract class OfflineToolUtils Collections.addAll(allowedThreadNames, EXTRA_JDK_THREADS); Collections.addAll(allowedThreadNames, optionalThreadNames); - if (allowNonDefaultMemtableThreads && DatabaseDescriptor.getMemtableConfigurations().containsKey("default")) + if (allowNonDefaultMemtableThreads + && (DatabaseDescriptor.getMemtableConfigurations() != null + && DatabaseDescriptor.getMemtableConfigurations().containsKey("default"))) Collections.addAll(allowedThreadNames, NON_DEFAULT_MEMTABLE_THREADS); var allowedRegexes = allowedThreadNames.stream() @@ -119,6 +121,7 @@ public abstract class OfflineToolUtils var badThreads = Arrays.stream(threads.getThreadInfo(threads.getAllThreadIds())) .filter(Objects::nonNull) .filter(threadInfo -> allowedRegexes.stream().noneMatch(pattern -> pattern.matcher(threadInfo.getThreadName()).matches())) + .filter(threadInfo -> !allowedThreadNames.contains(threadInfo.getThreadName())) .collect(Collectors.toSet()); if (!badThreads.isEmpty()) diff --git a/tools/bin/addtocmstool b/tools/bin/cmsofflinetool similarity index 93% rename from tools/bin/addtocmstool rename to tools/bin/cmsofflinetool index 3721639b6f..96ada2f195 100755 --- a/tools/bin/addtocmstool +++ b/tools/bin/cmsofflinetool @@ -43,7 +43,7 @@ fi "$JAVA" $JAVA_AGENT -ea -cp "$CLASSPATH" $JVM_OPTS -Xmx$MAX_HEAP_SIZE \ -Dcassandra.storagedir="$cassandra_storagedir" \ - -Dlog4j.configurationFile=log4j2-tools.xml \ - org.apache.cassandra.tools.TransformClusterMetadataHelper "$@" + -Dlogback.configurationFile=logback-tools.xml \ + org.apache.cassandra.tools.CMSOfflineTool "$@" # vi:ai sw=4 ts=4 tw=0 et