From 82acd3e9503452998988c0dc661f691ef94f7b4d Mon Sep 17 00:00:00 2001 From: Ariel Weisberg Date: Fri, 18 Aug 2023 16:48:39 -0400 Subject: [PATCH] Allow exceptions to be propagated remotely https://github.com/apache/cassandra-accord/pull/56 Patch by Ariel Weisberg; Reviewed by David Capwell for CASSANDRA-18779 --- modules/accord | 2 +- .../exceptions/ExceptionSerializer.java | 222 ++++++++++++++++++ .../cassandra/exceptions/RequestFailure.java | 164 +++++++++++++ .../exceptions/RequestFailureReason.java | 7 + .../cassandra/hints/HintsDispatcher.java | 17 +- .../apache/cassandra/locator/ReplicaPlan.java | 25 +- .../org/apache/cassandra/net/InboundSink.java | 10 +- .../org/apache/cassandra/net/Message.java | 12 +- .../apache/cassandra/net/MessageDelivery.java | 20 +- .../cassandra/net/MessagingService.java | 17 +- .../apache/cassandra/net/RequestCallback.java | 3 +- .../net/RequestCallbackWithFailure.java | 4 +- .../cassandra/net/RequestCallbacks.java | 8 +- .../cassandra/net/ResponseVerbHandler.java | 4 +- src/java/org/apache/cassandra/net/Verb.java | 71 ++++-- .../apache/cassandra/repair/SnapshotTask.java | 8 +- .../repair/messages/RepairMessage.java | 13 +- .../service/AbstractWriteResponseHandler.java | 5 +- .../service/ActiveRepairService.java | 7 +- .../service/BatchlogResponseHandler.java | 6 +- .../service/FailureRecordingCallback.java | 5 +- .../cassandra/service/StorageProxy.java | 14 +- .../service/TruncateResponseHandler.java | 10 +- .../service/accord/AccordCallback.java | 15 +- .../service/accord/AccordMessageSink.java | 20 +- .../service/accord/AccordSyncPropagator.java | 4 +- .../apache/cassandra/service/paxos/Paxos.java | 19 +- .../cassandra/service/paxos/PaxosCommit.java | 9 +- .../cassandra/service/paxos/PaxosPrepare.java | 48 +++- .../service/paxos/PaxosPrepareRefresh.java | 13 +- .../cassandra/service/paxos/PaxosPropose.java | 6 +- .../cassandra/service/paxos/PaxosRepair.java | 4 +- .../service/paxos/PaxosRequestCallback.java | 8 +- .../paxos/cleanup/PaxosCleanupComplete.java | 11 +- .../paxos/cleanup/PaxosCleanupSession.java | 8 +- .../cleanup/PaxosFinishPrepareCleanup.java | 8 +- .../cleanup/PaxosStartPrepareCleanup.java | 14 +- .../cassandra/service/reads/ReadCallback.java | 7 +- .../reads/thresholds/WarningContext.java | 14 +- .../cassandra/tcm/PaxosBackedProcessor.java | 8 +- .../apache/cassandra/tcm/RemoteProcessor.java | 5 +- .../cassandra/tcm/migration/Election.java | 1 - .../tcm/sequences/ProgressBarrier.java | 8 +- .../apache/cassandra/utils/FBUtilities.java | 3 - test/data/config/version=5.0-alpha1.yml | 2 +- .../simulator/systems/SimulatedAction.java | 4 +- .../config/ConfigCompatibilityTest.java | 1 + ...nterMutationVerbHandlerOutOfRangeTest.java | 14 +- .../db/MutationVerbHandlerOutOfRangeTest.java | 14 +- .../ReadCommandVerbHandlerOutOfRangeTest.java | 14 +- .../exceptions/RemoteExceptionTest.java | 112 +++++++++ .../apache/cassandra/net/ConnectionTest.java | 6 +- .../cassandra/net/MessageDeliveryTest.java | 4 +- .../org/apache/cassandra/net/MessageTest.java | 78 ++++-- .../net/SimulatedMessageDelivery.java | 12 +- .../apache/cassandra/repair/FuzzTestBase.java | 10 +- .../repair/messages/RepairMessageTest.java | 10 +- .../service/WriteResponseHandlerTest.java | 14 +- .../service/accord/AccordJournalTest.java | 9 + .../accord/AccordSyncPropagatorTest.java | 8 +- .../paxos/PaxosVerbHandlerOutOfRangeTest.java | 12 +- .../service/reads/ReadExecutorTest.java | 19 +- .../tcm/DiscoverySimulationTest.java | 4 +- .../tcm/sequences/ProgressBarrierTest.java | 4 +- 64 files changed, 960 insertions(+), 278 deletions(-) create mode 100644 src/java/org/apache/cassandra/exceptions/ExceptionSerializer.java create mode 100644 src/java/org/apache/cassandra/exceptions/RequestFailure.java create mode 100644 test/unit/org/apache/cassandra/exceptions/RemoteExceptionTest.java diff --git a/modules/accord b/modules/accord index 2ad55e03c4..91336705bd 160000 --- a/modules/accord +++ b/modules/accord @@ -1 +1 @@ -Subproject commit 2ad55e03c43ce074cdf5e36cfa14cb4278c2dc0f +Subproject commit 91336705bde8332954e849219d73205d68fa168a diff --git a/src/java/org/apache/cassandra/exceptions/ExceptionSerializer.java b/src/java/org/apache/cassandra/exceptions/ExceptionSerializer.java new file mode 100644 index 0000000000..de379739a3 --- /dev/null +++ b/src/java/org/apache/cassandra/exceptions/ExceptionSerializer.java @@ -0,0 +1,222 @@ +/* + * 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.exceptions; + +import java.io.IOException; +import java.util.HashMap; +import java.util.IdentityHashMap; +import java.util.Map; +import java.util.Set; + +import org.apache.cassandra.io.IVersionedSerializer; +import org.apache.cassandra.io.util.DataInputPlus; +import org.apache.cassandra.io.util.DataOutputPlus; +import org.apache.cassandra.utils.ArraySerializers; +import org.apache.cassandra.utils.FBUtilities; +import org.apache.cassandra.utils.NullableSerializer; + +import static java.util.Collections.newSetFromMap; +import static org.apache.cassandra.db.TypeSizes.sizeof; +import static org.apache.cassandra.db.TypeSizes.sizeofUnsignedVInt; + +/** + * Support for serializing exceptions without a dependency on being able to instantiate the exception class + * on the other side to eliminate dependencies across versions. + * + * This is still slightly more flexible than sending a string representation of the exception because it's still an exception so using it + * as a cause or suppressed exception works and it is formatted nicely as if it were another local exception. + */ +public class ExceptionSerializer +{ + public static class RemoteException extends RuntimeException + { + private final String originalClass; + + public RemoteException(String originalClass, String originalMessage, StackTraceElement[] stackTrace) + { + super(originalMessage); + this.originalClass = originalClass; + setStackTrace(stackTrace); + } + + private void initSuppressedAndCause(RemoteException cause, RemoteException[] suppressed) + { + initCause(cause); + for (RemoteException e : suppressed) + addSuppressed(e); + } + + @Override + public String toString() + { + String message = getMessage(); + return message != null ? originalClass + ": " + message : originalClass; + } + } + + static String getMessageWithOriginatingHost(Throwable t, boolean isFirstException) + { + if (isFirstException) + return "Remote exception from host " + FBUtilities.getBroadcastAddressAndPort().toString() + (t.getLocalizedMessage() != null ? " - " + t.getLocalizedMessage() : ""); + else + return t.getLocalizedMessage(); + } + + private static final IVersionedSerializer stackTraceElementSerializer = new IVersionedSerializer() + { + @Override + public void serialize(StackTraceElement t, DataOutputPlus out, int version) throws IOException + { + out.writeUTF(t.getClassName()); + out.writeUTF(t.getMethodName()); + out.writeBoolean(t.getFileName() != null); + if (t.getFileName() != null) + out.writeUTF(t.getFileName()); + out.writeUnsignedVInt32(t.getLineNumber()); + } + + @Override + public StackTraceElement deserialize(DataInputPlus in, int version) throws IOException + { + String className = in.readUTF(); + String methodName = in.readUTF(); + String fileName = null; + if (in.readBoolean()) + fileName = in.readUTF(); + int lineNumber = in.readUnsignedVInt32(); + return new StackTraceElement(className, methodName, fileName, lineNumber); + } + + @Override + public long serializedSize(StackTraceElement t, int version) + { + long size = sizeof(t.getClassName()) + + sizeof(t.getMethodName()) + + sizeof(t.getFileName() != null) + + sizeofUnsignedVInt(t.getLineNumber()); + if (t.getFileName() != null) + size += sizeof(t.getFileName()); + return size; + } + }; + + public static final IVersionedSerializer remoteExceptionSerializer = new IVersionedSerializer() + { + @Override + public void serialize(Throwable t, DataOutputPlus out, int version) throws IOException + { + Map alreadySerialized = new IdentityHashMap<>(); + serializeNextException(t, out, true, version, 0, alreadySerialized); + } + + private int serializeNextException(Throwable t, DataOutputPlus out, boolean isFirstException, int version, int nextExceptionId, Map alreadySerialized) throws IOException + { + if (alreadySerialized.containsKey(t)) + { + out.writeInt(alreadySerialized.get(t)); + return nextExceptionId; + } + else + { + alreadySerialized.put(t, nextExceptionId); + out.writeInt(nextExceptionId); + nextExceptionId++; + } + + out.writeUTF(t.getClass().getName()); + String message = getMessageWithOriginatingHost(t, isFirstException); + out.writeBoolean(message != null); + if (message != null) + out.writeUTF(message); + ArraySerializers.serializeArray(t.getStackTrace(), out, version, stackTraceElementSerializer); + + // Do cause and suppressed last so they can reference back to previously partially deserialized exceptions + out.writeBoolean(t.getCause() != null); + if (t.getCause() != null) + nextExceptionId = serializeNextException(t.getCause(), out, false, version, nextExceptionId, alreadySerialized); + out.writeUnsignedVInt32(t.getSuppressed().length); + for (Throwable suppressed : t.getSuppressed()) + nextExceptionId = serializeNextException(suppressed, out, false, version, nextExceptionId, alreadySerialized); + + return nextExceptionId; + } + + @Override + public Throwable deserialize(DataInputPlus in, int version) throws IOException + { + Map alreadyDeserialized = new HashMap<>(); + return deserializeNextException(in, version, alreadyDeserialized); + } + + private Throwable deserializeNextException(DataInputPlus in, int version, Map alreadyDeserialized) throws IOException + { + int nextExceptionId = in.readInt(); + Throwable alreadyDeserializedThrowable = alreadyDeserialized.get(nextExceptionId); + if (alreadyDeserializedThrowable != null) + return alreadyDeserializedThrowable; + + String originalClass = in.readUTF(); + String originalMessage = null; + if (in.readBoolean()) + originalMessage = in.readUTF(); + + StackTraceElement[] stackTrace = ArraySerializers.deserializeArray(in, version, stackTraceElementSerializer, size -> new StackTraceElement[size]); + RemoteException deserializedException = new RemoteException(originalClass, originalMessage, stackTrace); + deserializedException.setStackTrace(stackTrace); + alreadyDeserialized.put(nextExceptionId, deserializedException); + + // Do cause and suppressed last after alreadyDeserialized contains the exception we just processsed + RemoteException cause = in.readBoolean() ? (RemoteException)deserializeNextException(in, version, alreadyDeserialized) : null; + RemoteException[] suppressed = new RemoteException[in.readUnsignedVInt32()]; + for (int i = 0; i < suppressed.length; i++) + suppressed[i] = (RemoteException)deserializeNextException(in, version, alreadyDeserialized); + deserializedException.initSuppressedAndCause(cause, suppressed); + + return deserializedException; + } + + @Override + public long serializedSize(Throwable t, int version) + { + Set alreadySeen = newSetFromMap(new IdentityHashMap<>()); + return nextExceptionSerializedSize(t, version, true, alreadySeen); + } + + private long nextExceptionSerializedSize(Throwable t, int version, boolean isFirstException, Set alreadySeen) + { + if (!alreadySeen.add(t)) + return sizeof(42); // Exception ID from the last time it was serialized + + String message = getMessageWithOriginatingHost(t, isFirstException); + long size = sizeof(42) + // Exception ID generated during serialization + sizeof(t.getClass().getName()) + + sizeof(message != null) + + (message != null ? sizeof(message) : 0) + + sizeof(t.getCause() != null) + + (t.getCause() != null ? nextExceptionSerializedSize(t.getCause(), version, false, alreadySeen) : 0) + + sizeofUnsignedVInt(t.getSuppressed().length); + size += ArraySerializers.serializedArraySize(t.getStackTrace(), version, stackTraceElementSerializer); + for (Throwable suppressed : t.getSuppressed()) + size += nextExceptionSerializedSize(suppressed, version, false, alreadySeen); + return size; + } + }; + + public static final IVersionedSerializer nullableRemoteExceptionSerializer = NullableSerializer.wrap(remoteExceptionSerializer); +} diff --git a/src/java/org/apache/cassandra/exceptions/RequestFailure.java b/src/java/org/apache/cassandra/exceptions/RequestFailure.java new file mode 100644 index 0000000000..39700d795c --- /dev/null +++ b/src/java/org/apache/cassandra/exceptions/RequestFailure.java @@ -0,0 +1,164 @@ +/* + * 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.exceptions; + +import java.io.IOException; +import javax.annotation.Nonnull; +import javax.annotation.Nullable; + +import org.apache.cassandra.db.filter.TombstoneOverwhelmingException; +import org.apache.cassandra.io.IVersionedSerializer; +import org.apache.cassandra.io.util.DataInputPlus; +import org.apache.cassandra.io.util.DataOutputPlus; +import org.apache.cassandra.net.MessagingService; +import org.apache.cassandra.tcm.NotCMSException; + +import static com.google.common.base.Preconditions.checkNotNull; +import static org.apache.cassandra.exceptions.ExceptionSerializer.nullableRemoteExceptionSerializer; + +/** + * Allow inclusion of a serialized exception in failure response messages + * This continues to use the same verb as the old failure response (whether a message payload or parameter) + * and has a nullable failure field that may contain a serialized in later versions. + */ +public class RequestFailure +{ + public static final RequestFailure UNKNOWN = new RequestFailure(RequestFailureReason.UNKNOWN); + public static final RequestFailure READ_TOO_MANY_TOMBSTONES = new RequestFailure(RequestFailureReason.READ_TOO_MANY_TOMBSTONES); + public static final RequestFailure READ_TOO_MANY_INDEXES = new RequestFailure(RequestFailureReason.READ_TOO_MANY_INDEXES); + public static final RequestFailure TIMEOUT = new RequestFailure(RequestFailureReason.TIMEOUT); + public static final RequestFailure INCOMPATIBLE_SCHEMA = new RequestFailure(RequestFailureReason.INCOMPATIBLE_SCHEMA); + public static final RequestFailure READ_SIZE = new RequestFailure(RequestFailureReason.READ_SIZE); + public static final RequestFailure NODE_DOWN = new RequestFailure(RequestFailureReason.NODE_DOWN); + public static final RequestFailure NOT_CMS = new RequestFailure(RequestFailureReason.NOT_CMS); + public static final RequestFailure INVALID_ROUTING = new RequestFailure(RequestFailureReason.INVALID_ROUTING); + public static final RequestFailure INDEX_NOT_AVAILABLE = new RequestFailure(RequestFailureReason.INDEX_NOT_AVAILABLE); + public static final RequestFailure COORDINATOR_BEHIND = new RequestFailure(RequestFailureReason.COORDINATOR_BEHIND); + + static + { + // Validate all reasons are handled + for (RequestFailureReason reason : RequestFailureReason.values()) + forReason(reason); + } + + // Allow RequestFailureReason to force class load to check failure reasons are handled + public static void init() {} + + public static final IVersionedSerializer serializer = new IVersionedSerializer() + { + @Override + public void serialize(RequestFailure t, DataOutputPlus out, int version) throws IOException + { + RequestFailureReason.serializer.serialize(t.reason, out, version); + if (version >= MessagingService.VERSION_50) + nullableRemoteExceptionSerializer.serialize(t.failure, out, version); + } + + @Override + public RequestFailure deserialize(DataInputPlus in, int version) throws IOException + { + RequestFailureReason reason = RequestFailureReason.serializer.deserialize(in, version); + Throwable failure = null; + if (version >= MessagingService.VERSION_50) + failure = nullableRemoteExceptionSerializer.deserialize(in, version); + if (failure == null) + return forReason(reason); + else + return new RequestFailure(reason, failure); + } + + @Override + public long serializedSize(RequestFailure t, int version) + { + long size = RequestFailureReason.serializer.serializedSize(t.reason, version); + if (version >= MessagingService.VERSION_50) + size += nullableRemoteExceptionSerializer.serializedSize(t.failure, version); + return size; + } + }; + + @Nonnull + public final RequestFailureReason reason; + + @Nullable + public final Throwable failure; + + public static RequestFailure forException(Throwable t) + { + if (t instanceof TombstoneOverwhelmingException) + return READ_TOO_MANY_TOMBSTONES; + + if (t instanceof IncompatibleSchemaException) + return INCOMPATIBLE_SCHEMA; + + if (t instanceof NotCMSException) + return NOT_CMS; + + if (t instanceof InvalidRoutingException) + return INVALID_ROUTING; + + return UNKNOWN; + } + + public static RequestFailure forReason(RequestFailureReason reason) + { + switch (reason) + { + default: throw new IllegalStateException("Unhandled request failure reason " + reason); + case UNKNOWN: return UNKNOWN; + case READ_TOO_MANY_TOMBSTONES: return READ_TOO_MANY_TOMBSTONES; + case READ_TOO_MANY_INDEXES: return READ_TOO_MANY_INDEXES; + case TIMEOUT: return TIMEOUT; + case INCOMPATIBLE_SCHEMA: return INCOMPATIBLE_SCHEMA; + case READ_SIZE: return READ_SIZE; + case NODE_DOWN: return NODE_DOWN; + case NOT_CMS: return NOT_CMS; + case INVALID_ROUTING: return INVALID_ROUTING; + case INDEX_NOT_AVAILABLE: return INDEX_NOT_AVAILABLE; + case COORDINATOR_BEHIND: return COORDINATOR_BEHIND; + } + } + + private RequestFailure(RequestFailureReason reason) + { + this(reason, null); + } + + public RequestFailure(@Nonnull Throwable failure) + { + this(RequestFailureReason.UNKNOWN, failure); + } + + public RequestFailure(@Nonnull RequestFailureReason reason, @Nullable Throwable failure) + { + checkNotNull(reason); + this.reason = reason; + this.failure = failure; + } + + @Override + public String toString() + { + return "RequestFailure{" + + "reason=" + reason + + ", failure='" + failure + '\'' + + '}'; + } +} diff --git a/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java b/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java index 9faff584f1..560b8d68e0 100644 --- a/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java +++ b/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java @@ -49,6 +49,13 @@ public enum RequestFailureReason COORDINATOR_BEHIND (10), // The following codes have been ported from an external fork, where they were offset explicitly to avoid conflicts. INDEX_BUILD_IN_PROGRESS (503); + + static + { + // Load RequestFailure class to check that all request failure reasons are handled + RequestFailure.init(); + } + public static final Serializer serializer = new Serializer(); public final int code; diff --git a/src/java/org/apache/cassandra/hints/HintsDispatcher.java b/src/java/org/apache/cassandra/hints/HintsDispatcher.java index b627338543..ce1f7282a6 100644 --- a/src/java/org/apache/cassandra/hints/HintsDispatcher.java +++ b/src/java/org/apache/cassandra/hints/HintsDispatcher.java @@ -18,7 +18,10 @@ package org.apache.cassandra.hints; import java.nio.ByteBuffer; -import java.util.*; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Iterator; +import java.util.UUID; import java.util.function.BooleanSupplier; import java.util.function.Function; @@ -26,17 +29,19 @@ import com.google.common.util.concurrent.RateLimiter; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.apache.cassandra.net.RequestCallback; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.io.util.File; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.metrics.HintsServiceMetrics; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessagingService; +import org.apache.cassandra.net.RequestCallback; import org.apache.cassandra.utils.concurrent.Condition; - -import static org.apache.cassandra.hints.HintsDispatcher.Callback.Outcome.*; +import static org.apache.cassandra.hints.HintsDispatcher.Callback.Outcome.FAILURE; +import static org.apache.cassandra.hints.HintsDispatcher.Callback.Outcome.INTERRUPTED; +import static org.apache.cassandra.hints.HintsDispatcher.Callback.Outcome.SUCCESS; +import static org.apache.cassandra.hints.HintsDispatcher.Callback.Outcome.TIMEOUT; import static org.apache.cassandra.metrics.HintsServiceMetrics.updateDelayMetrics; import static org.apache.cassandra.net.Verb.HINT_REQ; import static org.apache.cassandra.utils.MonotonicClock.Global.approxTime; @@ -246,7 +251,7 @@ final class HintsDispatcher implements AutoCloseable } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failureMessage) { outcome = FAILURE; condition.signalAll(); diff --git a/src/java/org/apache/cassandra/locator/ReplicaPlan.java b/src/java/org/apache/cassandra/locator/ReplicaPlan.java index 7d08b341b8..ee8198aac3 100644 --- a/src/java/org/apache/cassandra/locator/ReplicaPlan.java +++ b/src/java/org/apache/cassandra/locator/ReplicaPlan.java @@ -18,16 +18,6 @@ package org.apache.cassandra.locator; -import com.google.common.collect.Iterables; -import org.apache.cassandra.db.ConsistencyLevel; -import org.apache.cassandra.db.Keyspace; -import org.apache.cassandra.db.PartitionPosition; -import org.apache.cassandra.dht.AbstractBounds; -import org.apache.cassandra.dht.Token; -import org.apache.cassandra.tcm.Epoch; -import org.apache.cassandra.exceptions.RequestFailureReason; -import org.apache.cassandra.tcm.ClusterMetadata; - import java.util.List; import java.util.concurrent.CopyOnWriteArrayList; import java.util.function.BiFunction; @@ -35,6 +25,17 @@ import java.util.function.Function; import java.util.function.Predicate; import java.util.function.Supplier; +import com.google.common.collect.Iterables; + +import org.apache.cassandra.db.ConsistencyLevel; +import org.apache.cassandra.db.Keyspace; +import org.apache.cassandra.db.PartitionPosition; +import org.apache.cassandra.dht.AbstractBounds; +import org.apache.cassandra.dht.Token; +import org.apache.cassandra.exceptions.RequestFailure; +import org.apache.cassandra.tcm.ClusterMetadata; +import org.apache.cassandra.tcm.Epoch; + public interface ReplicaPlan, P extends ReplicaPlan> { Epoch epoch(); @@ -49,7 +50,7 @@ public interface ReplicaPlan, P extends ReplicaPlan P withContacts(E contacts); void collectSuccess(InetAddressAndPort inetAddressAndPort); - void collectFailure(InetAddressAndPort inetAddressAndPort, RequestFailureReason t); + void collectFailure(InetAddressAndPort inetAddressAndPort, RequestFailure t); boolean stillAppliesTo(ClusterMetadata newMetadata); interface ForRead, P extends ReplicaPlan.ForRead> extends ReplicaPlan @@ -115,7 +116,7 @@ public interface ReplicaPlan, P extends ReplicaPlan contacted.add(addr); } - public void collectFailure(InetAddressAndPort inetAddressAndPort, RequestFailureReason t) {} + public void collectFailure(InetAddressAndPort inetAddressAndPort, RequestFailure t) {} } diff --git a/src/java/org/apache/cassandra/net/InboundSink.java b/src/java/org/apache/cassandra/net/InboundSink.java index 2e8c8413dc..7fbc50a201 100644 --- a/src/java/org/apache/cassandra/net/InboundSink.java +++ b/src/java/org/apache/cassandra/net/InboundSink.java @@ -30,11 +30,11 @@ import net.openhft.chronicle.core.util.ThrowingConsumer; import org.apache.cassandra.db.filter.TombstoneOverwhelmingException; import org.apache.cassandra.exceptions.CoordinatorBehindException; import org.apache.cassandra.exceptions.InvalidRoutingException; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.index.IndexNotAvailableException; import org.apache.cassandra.locator.InetAddressAndPort; -import org.apache.cassandra.tcm.Epoch; import org.apache.cassandra.tcm.ClusterMetadata; +import org.apache.cassandra.tcm.Epoch; import org.apache.cassandra.tcm.NotCMSException; import org.apache.cassandra.utils.NoSpamLogger; @@ -109,9 +109,9 @@ public class InboundSink implements InboundMessageHandlers.MessageConsumer if (header.callBackOnFailure()) { InetAddressAndPort to = header.respondTo() != null ? header.respondTo() : header.from; - Message response = Message.failureResponse(header.id, - header.expiresAtNanos, - RequestFailureReason.forException(failure)); + Message response = Message.failureResponse(header.id, + header.expiresAtNanos, + RequestFailure.forException(failure)); messaging.send(response, to); } } diff --git a/src/java/org/apache/cassandra/net/Message.java b/src/java/org/apache/cassandra/net/Message.java index f6ca39cd85..146f8c1119 100644 --- a/src/java/org/apache/cassandra/net/Message.java +++ b/src/java/org/apache/cassandra/net/Message.java @@ -36,6 +36,7 @@ import org.slf4j.LoggerFactory; import accord.messages.ReplyContext; import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.io.IVersionedAsymmetricSerializer; import org.apache.cassandra.io.IVersionedSerializer; @@ -353,12 +354,17 @@ public class Message implements ReplyContext } /** Builds a failure response Message with an explicit reason, and fields inferred from request Message */ - public Message failureResponse(RequestFailureReason reason) + public Message failureResponse(RequestFailureReason reason) { - return failureResponse(id(), expiresAtNanos(), reason); + return failureResponse(reason, null); } - static Message failureResponse(long id, long expiresAtNanos, RequestFailureReason reason) + public Message failureResponse(RequestFailureReason reason, @Nullable Throwable failure) + { + return failureResponse(id(), expiresAtNanos(), new RequestFailure(reason, failure)); + } + + static Message failureResponse(long id, long expiresAtNanos, RequestFailure reason) { return outWithParam(id, Verb.FAILURE_RSP, expiresAtNanos, reason, null, null); } diff --git a/src/java/org/apache/cassandra/net/MessageDelivery.java b/src/java/org/apache/cassandra/net/MessageDelivery.java index 0d052cb3d8..7c36d73c14 100644 --- a/src/java/org/apache/cassandra/net/MessageDelivery.java +++ b/src/java/org/apache/cassandra/net/MessageDelivery.java @@ -29,6 +29,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.utils.Backoff; @@ -63,7 +64,7 @@ public interface MessageDelivery } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason reason) + public void onFailure(InetAddressAndPort from, RequestFailure reason) { logger.info("Received failure in response to {} from {}: {}", verb, from, reason); cdl.decrement(); @@ -110,6 +111,11 @@ public interface MessageDelivery } public void respond(V response, Message message); public default void respondWithFailure(RequestFailureReason reason, Message message) + { + respondWithFailure(RequestFailure.forReason(reason), message); + } + + public default void respondWithFailure(RequestFailure reason, Message message) { send(Message.failureResponse(message.id(), message.expiresAtNanos(), reason), message.respondTo()); } @@ -121,12 +127,12 @@ public interface MessageDelivery interface RetryPredicate { - boolean test(int attempt, InetAddressAndPort from, RequestFailureReason failure); + boolean test(int attempt, InetAddressAndPort from, RequestFailure failure); } interface RetryErrorMessage { - String apply(int attempt, ResponseFailureReason retryFailure, @Nullable InetAddressAndPort from, @Nullable RequestFailureReason reason); + String apply(int attempt, ResponseFailureReason retryFailure, @Nullable InetAddressAndPort from, @Nullable RequestFailure reason); } private static void sendWithRetries(MessageDelivery messaging, @@ -157,7 +163,7 @@ public interface MessageDelivery } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failure) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { if (!backoff.mayRetry(attempt)) { @@ -212,11 +218,11 @@ public interface MessageDelivery class FailedResponseException extends IllegalStateException { public final InetAddressAndPort from; - public final RequestFailureReason failure; + public final RequestFailure failure; - public FailedResponseException(InetAddressAndPort from, RequestFailureReason failure, String message) + public FailedResponseException(InetAddressAndPort from, RequestFailure failure, String message) { - super(message); + super(message, failure.failure); this.from = from; this.failure = failure; } diff --git a/src/java/org/apache/cassandra/net/MessagingService.java b/src/java/org/apache/cassandra/net/MessagingService.java index 9215b7c975..12d3f17cd4 100644 --- a/src/java/org/apache/cassandra/net/MessagingService.java +++ b/src/java/org/apache/cassandra/net/MessagingService.java @@ -39,6 +39,7 @@ import org.apache.cassandra.concurrent.ScheduledExecutors; import org.apache.cassandra.concurrent.Stage; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.SystemKeyspace; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.locator.Replica; @@ -71,7 +72,7 @@ import static org.apache.cassandra.utils.Throwables.maybeFail; * message is received, {@link RequestCallback#onResponse(Message)} method will be invoked on the * provided callback - in case of a success response. In case of a failure response (see {@link Verb#FAILURE_RSP}), * or if a response doesn't arrive within verb's configured expiry time, - * {@link RequestCallback#onFailure(InetAddressAndPort, RequestFailureReason)} will be invoked instead. + * {@link RequestCallback#onFailure(InetAddressAndPort, RequestFailure)} will be invoked instead. * 2. To send a response back, or a message that expects no response, use {@link #send(Message, InetAddressAndPort)} * method. * @@ -381,9 +382,9 @@ public class MessagingService extends MessagingServiceMBeanImpl implements Messa } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { - promise.tryFailure(new FailureResponseException(from, failureReason)); + promise.tryFailure(new FailureResponseException(from, failure)); } @Override @@ -400,11 +401,11 @@ public class MessagingService extends MessagingServiceMBeanImpl implements Messa private final InetAddressAndPort from; private final RequestFailureReason failureReason; - public FailureResponseException(InetAddressAndPort from, RequestFailureReason failureReason) + public FailureResponseException(InetAddressAndPort from, RequestFailure failureReason) { - super(String.format("Failure from %s: %s", from, failureReason.name())); + super(String.format("Failure from %s: %s", from, failureReason.reason.name()), failureReason.failure); this.from = from; - this.failureReason = failureReason; + this.failureReason = failureReason.reason; } public InetAddressAndPort from() @@ -499,7 +500,7 @@ public class MessagingService extends MessagingServiceMBeanImpl implements Messa } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failureReason) { future.setFailure(new RuntimeException(failureReason.toString())); } @@ -510,7 +511,7 @@ public class MessagingService extends MessagingServiceMBeanImpl implements Messa public void respondWithFailure(RequestFailureReason reason, Message message) { - Message r = Message.failureResponse(message.id(), message.expiresAtNanos(), reason); + Message r = Message.failureResponse(message.id(), message.expiresAtNanos(), new RequestFailure(reason, null)); if (r.header.hasFlag(MessageFlag.URGENT)) r = r.withFlag(MessageFlag.URGENT); send(r, message.respondTo()); diff --git a/src/java/org/apache/cassandra/net/RequestCallback.java b/src/java/org/apache/cassandra/net/RequestCallback.java index 14e0169b85..1265b1ea6c 100644 --- a/src/java/org/apache/cassandra/net/RequestCallback.java +++ b/src/java/org/apache/cassandra/net/RequestCallback.java @@ -19,6 +19,7 @@ package org.apache.cassandra.net; import java.util.Map; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.locator.InetAddressAndPort; @@ -38,7 +39,7 @@ public interface RequestCallback /** * Called when there is an exception on the remote node or timeout happens */ - default void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + default void onFailure(InetAddressAndPort from, RequestFailure failure) { } diff --git a/src/java/org/apache/cassandra/net/RequestCallbackWithFailure.java b/src/java/org/apache/cassandra/net/RequestCallbackWithFailure.java index 685797abeb..a7d807380e 100644 --- a/src/java/org/apache/cassandra/net/RequestCallbackWithFailure.java +++ b/src/java/org/apache/cassandra/net/RequestCallbackWithFailure.java @@ -18,7 +18,7 @@ package org.apache.cassandra.net; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.locator.InetAddressAndPort; public interface RequestCallbackWithFailure extends RequestCallback @@ -26,7 +26,7 @@ public interface RequestCallbackWithFailure extends RequestCallback /** * Called when there is an exception on the remote node or timeout happens */ - void onFailure(InetAddressAndPort from, RequestFailureReason failureReason); + void onFailure(InetAddressAndPort from, RequestFailure failure); /** * @return true if the callback should be invoked on failure diff --git a/src/java/org/apache/cassandra/net/RequestCallbacks.java b/src/java/org/apache/cassandra/net/RequestCallbacks.java index ee63c5a3e6..485efc30b1 100644 --- a/src/java/org/apache/cassandra/net/RequestCallbacks.java +++ b/src/java/org/apache/cassandra/net/RequestCallbacks.java @@ -21,17 +21,15 @@ import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.TimeoutException; - import javax.annotation.Nullable; import com.google.common.annotations.VisibleForTesting; - -import org.apache.cassandra.concurrent.ScheduledExecutorPlus; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.apache.cassandra.concurrent.ScheduledExecutorPlus; import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.locator.Replica; import org.apache.cassandra.metrics.InternodeOutboundMetrics; @@ -154,7 +152,7 @@ public class RequestCallbacks implements OutboundMessageCallbacks messagingService.markExpiredCallback(info.peer); if (info.invokeOnFailure()) - INTERNAL_RESPONSE.submit(() -> info.callback.onFailure(info.peer, RequestFailureReason.TIMEOUT)); + INTERNAL_RESPONSE.submit(() -> info.callback.onFailure(info.peer, RequestFailure.TIMEOUT)); } void shutdownNow(boolean expireCallbacks) diff --git a/src/java/org/apache/cassandra/net/ResponseVerbHandler.java b/src/java/org/apache/cassandra/net/ResponseVerbHandler.java index 36e5cf0670..6cecd2a415 100644 --- a/src/java/org/apache/cassandra/net/ResponseVerbHandler.java +++ b/src/java/org/apache/cassandra/net/ResponseVerbHandler.java @@ -24,7 +24,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.cassandra.concurrent.Stage; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.tcm.ClusterMetadataService; import org.apache.cassandra.tracing.Tracing; @@ -74,7 +74,7 @@ class ResponseVerbHandler implements IVerbHandler RequestCallback cb = callbackInfo.callback; if (message.isFailureResponse()) { - cb.onFailure(message.from(), (RequestFailureReason) message.payload); + cb.onFailure(message.from(), (RequestFailure) message.payload); } else { diff --git a/src/java/org/apache/cassandra/net/Verb.java b/src/java/org/apache/cassandra/net/Verb.java index dc8e6d4157..c87fa0f163 100644 --- a/src/java/org/apache/cassandra/net/Verb.java +++ b/src/java/org/apache/cassandra/net/Verb.java @@ -41,13 +41,13 @@ import org.apache.cassandra.db.ReadCommandVerbHandler; import org.apache.cassandra.db.ReadRepairVerbHandler; import org.apache.cassandra.db.ReadResponse; import org.apache.cassandra.db.SnapshotCommand; +import org.apache.cassandra.db.TruncateRequest; import org.apache.cassandra.db.TruncateResponse; import org.apache.cassandra.db.TruncateVerbHandler; -import org.apache.cassandra.db.TruncateRequest; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; +import org.apache.cassandra.gms.GossipDigestAck; import org.apache.cassandra.gms.GossipDigestAck2; import org.apache.cassandra.gms.GossipDigestAck2VerbHandler; -import org.apache.cassandra.gms.GossipDigestAck; import org.apache.cassandra.gms.GossipDigestAckVerbHandler; import org.apache.cassandra.gms.GossipDigestSyn; import org.apache.cassandra.gms.GossipDigestSynVerbHandler; @@ -68,17 +68,19 @@ import org.apache.cassandra.repair.messages.PrepareMessage; import org.apache.cassandra.repair.messages.SnapshotMessage; import org.apache.cassandra.repair.messages.StatusRequest; import org.apache.cassandra.repair.messages.StatusResponse; -import org.apache.cassandra.repair.messages.SyncResponse; import org.apache.cassandra.repair.messages.SyncRequest; -import org.apache.cassandra.repair.messages.ValidationResponse; +import org.apache.cassandra.repair.messages.SyncResponse; import org.apache.cassandra.repair.messages.ValidationRequest; +import org.apache.cassandra.repair.messages.ValidationResponse; import org.apache.cassandra.schema.SchemaMutationsSerializer; import org.apache.cassandra.schema.SchemaPullVerbHandler; import org.apache.cassandra.schema.SchemaPushVerbHandler; import org.apache.cassandra.schema.SchemaVersionVerbHandler; +import org.apache.cassandra.service.EchoVerbHandler; +import org.apache.cassandra.service.SnapshotVerbHandler; import org.apache.cassandra.service.accord.AccordService; -import org.apache.cassandra.service.accord.AccordSyncPropagator.Notification; import org.apache.cassandra.service.accord.AccordSyncPropagator; +import org.apache.cassandra.service.accord.AccordSyncPropagator.Notification; import org.apache.cassandra.service.accord.serializers.AcceptSerializers; import org.apache.cassandra.service.accord.serializers.ApplySerializers; import org.apache.cassandra.service.accord.serializers.BeginInvalidationSerializers; @@ -104,14 +106,18 @@ import org.apache.cassandra.service.paxos.PaxosPrepare; import org.apache.cassandra.service.paxos.PaxosPrepareRefresh; import org.apache.cassandra.service.paxos.PaxosPropose; import org.apache.cassandra.service.paxos.PaxosRepair; +import org.apache.cassandra.service.paxos.PrepareResponse; +import org.apache.cassandra.service.paxos.cleanup.PaxosCleanupComplete; import org.apache.cassandra.service.paxos.cleanup.PaxosCleanupHistory; import org.apache.cassandra.service.paxos.cleanup.PaxosCleanupRequest; import org.apache.cassandra.service.paxos.cleanup.PaxosCleanupResponse; -import org.apache.cassandra.service.paxos.cleanup.PaxosCleanupComplete; -import org.apache.cassandra.service.paxos.cleanup.PaxosStartPrepareCleanup; import org.apache.cassandra.service.paxos.cleanup.PaxosFinishPrepareCleanup; +import org.apache.cassandra.service.paxos.cleanup.PaxosStartPrepareCleanup; +import org.apache.cassandra.service.paxos.v1.PrepareVerbHandler; +import org.apache.cassandra.service.paxos.v1.ProposeVerbHandler; import org.apache.cassandra.streaming.DataMovement; import org.apache.cassandra.streaming.DataMovementVerbHandler; +import org.apache.cassandra.streaming.ReplicationDoneVerbHandler; import org.apache.cassandra.tcm.Discovery; import org.apache.cassandra.tcm.Epoch; import org.apache.cassandra.tcm.FetchCMSLog; @@ -122,26 +128,49 @@ import org.apache.cassandra.tcm.migration.CMSInitializationRequest; import org.apache.cassandra.tcm.sequences.DataMovements; import org.apache.cassandra.tcm.serialization.MessageSerializers; import org.apache.cassandra.utils.BooleanSerializer; -import org.apache.cassandra.service.EchoVerbHandler; -import org.apache.cassandra.service.SnapshotVerbHandler; -import org.apache.cassandra.service.paxos.PrepareResponse; -import org.apache.cassandra.service.paxos.v1.PrepareVerbHandler; -import org.apache.cassandra.service.paxos.v1.ProposeVerbHandler; -import org.apache.cassandra.streaming.ReplicationDoneVerbHandler; import org.apache.cassandra.utils.ReflectionUtils; import org.apache.cassandra.utils.TimeUUID; import org.apache.cassandra.utils.UUIDSerializer; import static java.util.concurrent.TimeUnit.NANOSECONDS; -import static org.apache.cassandra.concurrent.Stage.*; +import static org.apache.cassandra.concurrent.Stage.ANTI_ENTROPY; +import static org.apache.cassandra.concurrent.Stage.COUNTER_MUTATION; +import static org.apache.cassandra.concurrent.Stage.FETCH_LOG; +import static org.apache.cassandra.concurrent.Stage.GOSSIP; +import static org.apache.cassandra.concurrent.Stage.IMMEDIATE; +import static org.apache.cassandra.concurrent.Stage.INTERNAL_METADATA; +import static org.apache.cassandra.concurrent.Stage.INTERNAL_RESPONSE; +import static org.apache.cassandra.concurrent.Stage.MIGRATION; +import static org.apache.cassandra.concurrent.Stage.MISC; +import static org.apache.cassandra.concurrent.Stage.MUTATION; +import static org.apache.cassandra.concurrent.Stage.PAXOS_REPAIR; +import static org.apache.cassandra.concurrent.Stage.READ; +import static org.apache.cassandra.concurrent.Stage.REQUEST_RESPONSE; +import static org.apache.cassandra.concurrent.Stage.TRACING; import static org.apache.cassandra.net.ResponseHandlerSupplier.RESPONSE_HANDLER; -import static org.apache.cassandra.net.VerbTimeouts.*; -import static org.apache.cassandra.net.Verb.Kind.*; -import static org.apache.cassandra.net.Verb.Priority.*; +import static org.apache.cassandra.net.Verb.Kind.CUSTOM; +import static org.apache.cassandra.net.Verb.Kind.NORMAL; +import static org.apache.cassandra.net.Verb.Priority.P0; +import static org.apache.cassandra.net.Verb.Priority.P1; +import static org.apache.cassandra.net.Verb.Priority.P2; +import static org.apache.cassandra.net.Verb.Priority.P3; +import static org.apache.cassandra.net.Verb.Priority.P4; +import static org.apache.cassandra.net.VerbTimeouts.counterTimeout; +import static org.apache.cassandra.net.VerbTimeouts.longTimeout; +import static org.apache.cassandra.net.VerbTimeouts.noTimeout; +import static org.apache.cassandra.net.VerbTimeouts.pingTimeout; +import static org.apache.cassandra.net.VerbTimeouts.rangeTimeout; +import static org.apache.cassandra.net.VerbTimeouts.readTimeout; +import static org.apache.cassandra.net.VerbTimeouts.repairTimeout; +import static org.apache.cassandra.net.VerbTimeouts.repairValidationRspTimeout; +import static org.apache.cassandra.net.VerbTimeouts.repairWithBackoffTimeout; +import static org.apache.cassandra.net.VerbTimeouts.rpcTimeout; +import static org.apache.cassandra.net.VerbTimeouts.truncateTimeout; +import static org.apache.cassandra.net.VerbTimeouts.writeTimeout; import static org.apache.cassandra.tcm.ClusterMetadataService.commitRequestHandler; import static org.apache.cassandra.tcm.ClusterMetadataService.currentEpochRequestHandler; -import static org.apache.cassandra.tcm.ClusterMetadataService.logNotifyHandler; import static org.apache.cassandra.tcm.ClusterMetadataService.fetchLogRequestHandler; +import static org.apache.cassandra.tcm.ClusterMetadataService.logNotifyHandler; import static org.apache.cassandra.tcm.ClusterMetadataService.replicationHandler; /** @@ -306,7 +335,7 @@ public enum Verb ACCORD_SYNC_NOTIFY_REQ (151, P2, writeTimeout, IMMEDIATE, () -> Notification.listSerializer, () -> AccordSyncPropagator.verbHandler, ACCORD_SIMPLE_RSP ), // generic failure response - FAILURE_RSP (99, P0, noTimeout, REQUEST_RESPONSE, () -> RequestFailureReason.serializer, RESPONSE_HANDLER ), + FAILURE_RSP (99, P0, noTimeout, REQUEST_RESPONSE, () -> RequestFailure.serializer, RESPONSE_HANDLER ), // dummy verbs _TRACE (30, P1, rpcTimeout, TRACING, () -> NoPayload.serializer, () -> null ), @@ -438,7 +467,7 @@ public enum Verb // this is a little hacky, but reduces the number of parameters up top public boolean isResponse() { - return handler == RESPONSE_HANDLER; + return handler.get() == ResponseVerbHandler.instance; } @VisibleForTesting diff --git a/src/java/org/apache/cassandra/repair/SnapshotTask.java b/src/java/org/apache/cassandra/repair/SnapshotTask.java index ad45070cfb..a95e0668d8 100644 --- a/src/java/org/apache/cassandra/repair/SnapshotTask.java +++ b/src/java/org/apache/cassandra/repair/SnapshotTask.java @@ -19,10 +19,10 @@ package org.apache.cassandra.repair; import java.util.concurrent.RunnableFuture; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.locator.InetAddressAndPort; -import org.apache.cassandra.net.RequestCallback; import org.apache.cassandra.net.Message; +import org.apache.cassandra.net.RequestCallback; import org.apache.cassandra.repair.messages.RepairMessage; import org.apache.cassandra.repair.messages.SnapshotMessage; import org.apache.cassandra.utils.concurrent.AsyncFuture; @@ -81,9 +81,9 @@ public class SnapshotTask extends AsyncFuture implements Run } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { - task.tryFailure(new RuntimeException("Could not create snapshot at " + from + "; " + failureReason)); + task.tryFailure(new RuntimeException("Could not create snapshot at " + from + "; " + failure.reason)); } } } diff --git a/src/java/org/apache/cassandra/repair/messages/RepairMessage.java b/src/java/org/apache/cassandra/repair/messages/RepairMessage.java index 835f90fc68..e38a930bcc 100644 --- a/src/java/org/apache/cassandra/repair/messages/RepairMessage.java +++ b/src/java/org/apache/cassandra/repair/messages/RepairMessage.java @@ -36,6 +36,7 @@ import org.slf4j.LoggerFactory; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.config.RepairRetrySpec; import org.apache.cassandra.config.RetrySpec; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.metrics.RepairMetrics; import org.apache.cassandra.repair.SharedContext; import org.apache.cassandra.exceptions.RepairException; @@ -73,7 +74,7 @@ public abstract class RepairMessage } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failureReason) { } }; @@ -217,9 +218,9 @@ public abstract class RepairMessage finalCallback.onFailure(from, failure); return false; case RETRY: - if (failure == RequestFailureReason.TIMEOUT && allowRetry.get()) + if (failure.reason == RequestFailureReason.TIMEOUT && allowRetry.get()) return true; - maybeRecordRetry.accept(attempt, failure); + maybeRecordRetry.accept(attempt, failure.reason); finalCallback.onFailure(from, failure); return false; default: @@ -230,7 +231,7 @@ public abstract class RepairMessage switch (retryReason) { case MaxRetries: - maybeRecordRetry.accept(attempt, failure); + maybeRecordRetry.accept(attempt, failure.reason); finalCallback.onFailure(from, failure); return null; case Interrupted: @@ -255,9 +256,9 @@ public abstract class RepairMessage } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { - failureCallback.onFailure(RepairException.error(request.desc, PreviewKind.NONE, String.format("Got %s failure from %s: %s", verb, from, failureReason))); + failureCallback.onFailure(RepairException.error(request.desc, PreviewKind.NONE, String.format("Got %s failure from %s: %s", verb, from, failure.reason))); } @Override diff --git a/src/java/org/apache/cassandra/service/AbstractWriteResponseHandler.java b/src/java/org/apache/cassandra/service/AbstractWriteResponseHandler.java index 343bac1c4f..d085174043 100644 --- a/src/java/org/apache/cassandra/service/AbstractWriteResponseHandler.java +++ b/src/java/org/apache/cassandra/service/AbstractWriteResponseHandler.java @@ -36,6 +36,7 @@ import org.apache.cassandra.db.ConsistencyLevel; import org.apache.cassandra.db.IMutation; import org.apache.cassandra.db.Mutation; import org.apache.cassandra.db.WriteType; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.exceptions.WriteFailureException; import org.apache.cassandra.exceptions.WriteTimeoutException; @@ -295,7 +296,7 @@ public abstract class AbstractWriteResponseHandler implements RequestCallback } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { logger.trace("Got failure from {}", from); @@ -309,7 +310,7 @@ public abstract class AbstractWriteResponseHandler implements RequestCallback if (failureReasonByEndpoint == null) failureReasonByEndpoint = new ConcurrentHashMap<>(); } - failureReasonByEndpoint.put(from, failureReason); + failureReasonByEndpoint.put(from, failure.reason); logFailureOrTimeoutToIdealCLDelegate(); diff --git a/src/java/org/apache/cassandra/service/ActiveRepairService.java b/src/java/org/apache/cassandra/service/ActiveRepairService.java index 1a60bd3067..1592718bf1 100644 --- a/src/java/org/apache/cassandra/service/ActiveRepairService.java +++ b/src/java/org/apache/cassandra/service/ActiveRepairService.java @@ -64,6 +64,7 @@ import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; import org.apache.cassandra.exceptions.ConfigurationException; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.gms.ApplicationState; import org.apache.cassandra.gms.EndpointState; @@ -731,10 +732,10 @@ public class ActiveRepairService implements IEndpointStateChangeSubscriber, IFai } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { failedNodes.add(from.toString()); - if (failureReason == RequestFailureReason.TIMEOUT) + if (failure.reason == RequestFailureReason.TIMEOUT) { pending.set(-1); promise.setFailure(failRepairException(parentRepairSession, "Did not get replies from all endpoints.")); @@ -787,7 +788,7 @@ public class ActiveRepairService implements IEndpointStateChangeSubscriber, IFai } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { logger.debug("Failed to clean up parent repair session {} on {}. The uncleaned sessions will " + "be removed on a node restart. This should not be a problem unless you see thousands " + diff --git a/src/java/org/apache/cassandra/service/BatchlogResponseHandler.java b/src/java/org/apache/cassandra/service/BatchlogResponseHandler.java index 0fa2847700..dd2ebae915 100644 --- a/src/java/org/apache/cassandra/service/BatchlogResponseHandler.java +++ b/src/java/org/apache/cassandra/service/BatchlogResponseHandler.java @@ -20,7 +20,7 @@ package org.apache.cassandra.service; import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.WriteFailureException; import org.apache.cassandra.exceptions.WriteTimeoutException; import org.apache.cassandra.locator.InetAddressAndPort; @@ -55,9 +55,9 @@ public class BatchlogResponseHandler extends AbstractWriteResponseHandler cleanup.ackMutation(); } - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { - wrapped.onFailure(from, failureReason); + wrapped.onFailure(from, failure); } public boolean invokeOnFailure() diff --git a/src/java/org/apache/cassandra/service/FailureRecordingCallback.java b/src/java/org/apache/cassandra/service/FailureRecordingCallback.java index c4ca8e22f5..b672c4bf3f 100644 --- a/src/java/org/apache/cassandra/service/FailureRecordingCallback.java +++ b/src/java/org/apache/cassandra/service/FailureRecordingCallback.java @@ -25,6 +25,7 @@ import java.util.Map; import java.util.Set; import java.util.concurrent.atomic.AtomicReferenceFieldUpdater; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.net.RequestCallbackWithFailure; @@ -136,9 +137,9 @@ public abstract class FailureRecordingCallback implements RequestCallbackWith private static final AtomicReferenceFieldUpdater responsesUpdater = AtomicReferenceFieldUpdater.newUpdater(FailureRecordingCallback.class, FailureResponses.class, "failureResponses"); @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { - FailureResponses.push(responsesUpdater, this, from, failureReason); + FailureResponses.push(responsesUpdater, this, from, failure.reason); } protected void onFailureWithMutex(InetAddressAndPort from, RequestFailureReason failureReason) diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index c8eb3be597..e61c9187a3 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -91,7 +91,7 @@ import org.apache.cassandra.exceptions.ReadAbortException; import org.apache.cassandra.exceptions.ReadFailureException; import org.apache.cassandra.exceptions.ReadTimeoutException; import org.apache.cassandra.exceptions.RequestFailureException; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.RequestTimeoutException; import org.apache.cassandra.exceptions.UnavailableException; import org.apache.cassandra.exceptions.WriteFailureException; @@ -861,7 +861,7 @@ public class StorageProxy implements StorageProxyMBean { if (!(ex instanceof WriteTimeoutException)) logger.error("Failed to apply paxos commit locally : ", ex); - responseHandler.onFailure(FBUtilities.getBroadcastAddressAndPort(), RequestFailureReason.forException(ex)); + responseHandler.onFailure(FBUtilities.getBroadcastAddressAndPort(), RequestFailure.forException(ex)); } } @@ -1373,7 +1373,7 @@ public class StorageProxy implements StorageProxyMBean } catch (OverloadedException | WriteTimeoutException e) { - wrapper.handler.onFailure(FBUtilities.getBroadcastAddressAndPort(), RequestFailureReason.forException(e)); + wrapper.handler.onFailure(FBUtilities.getBroadcastAddressAndPort(), RequestFailure.forException(e)); } } } @@ -1723,7 +1723,7 @@ public class StorageProxy implements StorageProxyMBean { if (!(ex instanceof WriteTimeoutException)) logger.error("Failed to apply mutation locally : ", ex); - handler.onFailure(FBUtilities.getBroadcastAddressAndPort(), RequestFailureReason.forException(ex)); + handler.onFailure(FBUtilities.getBroadcastAddressAndPort(), RequestFailure.forException(ex)); } } @@ -2246,7 +2246,7 @@ public class StorageProxy implements StorageProxyMBean { // We track latency based on request processing time MessagingService.instance().metrics.recordSelfDroppedMessage(verb, MonotonicClock.Global.preciseTime.now() - requestTime.startedAtNanos(), NANOSECONDS); - handler.onFailure(FBUtilities.getBroadcastAddressAndPort(), RequestFailureReason.UNKNOWN); + handler.onFailure(FBUtilities.getBroadcastAddressAndPort(), RequestFailure.UNKNOWN); } if (!readRejected) @@ -2257,12 +2257,12 @@ public class StorageProxy implements StorageProxyMBean { if (t instanceof TombstoneOverwhelmingException) { - handler.onFailure(FBUtilities.getBroadcastAddressAndPort(), RequestFailureReason.READ_TOO_MANY_TOMBSTONES); + handler.onFailure(FBUtilities.getBroadcastAddressAndPort(), RequestFailure.READ_TOO_MANY_TOMBSTONES); logger.error(t.getMessage()); } else { - handler.onFailure(FBUtilities.getBroadcastAddressAndPort(), RequestFailureReason.UNKNOWN); + handler.onFailure(FBUtilities.getBroadcastAddressAndPort(), RequestFailure.UNKNOWN); throw t; } } diff --git a/src/java/org/apache/cassandra/service/TruncateResponseHandler.java b/src/java/org/apache/cassandra/service/TruncateResponseHandler.java index 54b1241006..566d9fa022 100644 --- a/src/java/org/apache/cassandra/service/TruncateResponseHandler.java +++ b/src/java/org/apache/cassandra/service/TruncateResponseHandler.java @@ -23,19 +23,19 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicInteger; -import org.apache.cassandra.utils.concurrent.Condition; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.cassandra.db.TruncateResponse; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.exceptions.TruncateException; import org.apache.cassandra.locator.InetAddressAndPort; -import org.apache.cassandra.net.RequestCallback; import org.apache.cassandra.net.Message; +import org.apache.cassandra.net.RequestCallback; +import org.apache.cassandra.utils.concurrent.Condition; import org.apache.cassandra.utils.concurrent.UncheckedInterruptedException; - import static java.util.concurrent.TimeUnit.NANOSECONDS; import static org.apache.cassandra.config.DatabaseDescriptor.getTruncateRpcTimeout; import static org.apache.cassandra.utils.Clock.Global.nanoTime; @@ -100,10 +100,10 @@ public class TruncateResponseHandler implements RequestCallback extends SafeCallback implements Request success(endpointMapper.mappedId(msg.from()), msg.payload); } - private static Throwable convertReason(RequestFailureReason reason) + private static Throwable convertFailureMessage(RequestFailure failure) { - return reason == RequestFailureReason.TIMEOUT ? + return failure.reason == RequestFailureReason.TIMEOUT ? new Timeout(null, null) : - new RuntimeException(reason.toString()); + new RuntimeException(failure.failure); } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { - logger.debug("Received failure {} from {} for {}", failureReason, from, this); + logger.debug("Received failure {} from {} for {}", failure, from, this); // TODO (now): we should distinguish timeout failures with some placeholder Exception - failure(endpointMapper.mappedId(from), convertReason(failureReason)); + failure(endpointMapper.mappedId(from), convertFailureMessage(failure)); } @Override diff --git a/src/java/org/apache/cassandra/service/accord/AccordMessageSink.java b/src/java/org/apache/cassandra/service/accord/AccordMessageSink.java index a6d8e4f341..1692e85fe8 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordMessageSink.java +++ b/src/java/org/apache/cassandra/service/accord/AccordMessageSink.java @@ -26,8 +26,6 @@ import java.util.Set; import com.google.common.base.Preconditions; import com.google.common.collect.ImmutableMap; - -import org.apache.cassandra.net.*; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -40,7 +38,12 @@ import accord.messages.MessageType; import accord.messages.Reply; import accord.messages.ReplyContext; import accord.messages.Request; +import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.locator.InetAddressAndPort; +import org.apache.cassandra.net.Message; +import org.apache.cassandra.net.MessageDelivery; +import org.apache.cassandra.net.MessagingService; +import org.apache.cassandra.net.Verb; public class AccordMessageSink implements MessageSink { @@ -93,6 +96,9 @@ public class AccordMessageSink implements MessageSink for (MessageType type : MessageType.values()) { + // Any request can receive a generic failure response + if (type == MessageType.FAILURE_RSP) + continue; if (!mapping.containsKey(type)) throw new AssertionError("Missing mapping for Accord MessageType " + type); } @@ -153,6 +159,16 @@ public class AccordMessageSink implements MessageSink messaging.send(replyMsg, endpoint); } + @Override + public void replyWithUnknownFailure(Node.Id replyingToNode, ReplyContext replyContext, Throwable failure) + { + Message replyTo = (Message) replyContext; + Message replyMsg = replyTo.failureResponse(RequestFailureReason.UNKNOWN, failure); + InetAddressAndPort endpoint = endpointMapper.mappedEndpoint(replyingToNode); + logger.debug("Replying with failure {} {} to {}", replyMsg.verb(), replyMsg.payload, endpoint); + messaging.send(replyMsg, endpoint); + } + private static void checkReplyType(Reply reply, Message replyTo) { Verb verb = getVerb(reply.type()); diff --git a/src/java/org/apache/cassandra/service/accord/AccordSyncPropagator.java b/src/java/org/apache/cassandra/service/accord/AccordSyncPropagator.java index 3e215a7e68..e16facee4a 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordSyncPropagator.java +++ b/src/java/org/apache/cassandra/service/accord/AccordSyncPropagator.java @@ -37,7 +37,7 @@ import org.agrona.collections.Int2ObjectHashMap; import org.agrona.collections.Long2ObjectHashMap; import org.apache.cassandra.concurrent.ScheduledExecutorPlus; import org.apache.cassandra.db.TypeSizes; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.gms.IFailureDetector; import org.apache.cassandra.io.IVersionedSerializer; import org.apache.cassandra.io.util.DataInputPlus; @@ -294,7 +294,7 @@ public class AccordSyncPropagator } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { scheduler.schedule(() -> AccordSyncPropagator.this.notify(to, notifications), 1, TimeUnit.MINUTES); } diff --git a/src/java/org/apache/cassandra/service/paxos/Paxos.java b/src/java/org/apache/cassandra/service/paxos/Paxos.java index a82396cd24..75392640a0 100644 --- a/src/java/org/apache/cassandra/service/paxos/Paxos.java +++ b/src/java/org/apache/cassandra/service/paxos/Paxos.java @@ -28,13 +28,11 @@ import java.util.concurrent.TimeUnit; import java.util.function.Function; import java.util.function.Predicate; import java.util.function.Supplier; - import javax.annotation.Nullable; import com.google.common.base.Preconditions; import com.google.common.collect.Iterators; import com.google.common.collect.Maps; - import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -73,6 +71,7 @@ import org.apache.cassandra.exceptions.ReadFailureException; import org.apache.cassandra.exceptions.ReadTimeoutException; import org.apache.cassandra.exceptions.RequestExecutionException; import org.apache.cassandra.exceptions.RequestFailureException; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.exceptions.RequestTimeoutException; import org.apache.cassandra.exceptions.UnavailableException; @@ -91,6 +90,8 @@ import org.apache.cassandra.service.FailureRecordingCallback.AsMap; import org.apache.cassandra.service.paxos.Commit.Proposal; import org.apache.cassandra.service.paxos.cleanup.PaxosRepairState; import org.apache.cassandra.tcm.ClusterMetadata; +import org.apache.cassandra.service.paxos.PaxosPrepare.FoundIncompleteAccepted; +import org.apache.cassandra.service.paxos.PaxosPrepare.FoundIncompleteCommitted; import org.apache.cassandra.service.reads.DataResolver; import org.apache.cassandra.service.reads.repair.NoopReadRepair; import org.apache.cassandra.tcm.Epoch; @@ -101,8 +102,6 @@ import org.apache.cassandra.transport.Dispatcher; import org.apache.cassandra.triggers.TriggerExecutor; import org.apache.cassandra.utils.CassandraVersion; import org.apache.cassandra.utils.FBUtilities; -import org.apache.cassandra.service.paxos.PaxosPrepare.FoundIncompleteAccepted; -import org.apache.cassandra.service.paxos.PaxosPrepare.FoundIncompleteCommitted; import org.apache.cassandra.utils.NoSpamLogger; import static java.util.Collections.emptyMap; @@ -111,10 +110,14 @@ import static java.util.concurrent.TimeUnit.SECONDS; import static org.apache.cassandra.config.CassandraRelevantProperties.PAXOS_LOG_TTL_LINEARIZABILITY_VIOLATIONS; import static org.apache.cassandra.config.CassandraRelevantProperties.PAXOS_MODERN_RELEASE; import static org.apache.cassandra.config.Config.PaxosVariant.v2_without_linearizable_reads_or_rejected_writes; +import static org.apache.cassandra.config.DatabaseDescriptor.getCasContentionTimeout; +import static org.apache.cassandra.config.DatabaseDescriptor.getWriteRpcTimeout; +import static org.apache.cassandra.db.ConsistencyLevel.LOCAL_QUORUM; +import static org.apache.cassandra.db.ConsistencyLevel.LOCAL_SERIAL; +import static org.apache.cassandra.db.ConsistencyLevel.QUORUM; +import static org.apache.cassandra.db.ConsistencyLevel.SERIAL; import static org.apache.cassandra.db.Keyspace.openAndGetStore; import static org.apache.cassandra.exceptions.RequestFailureReason.TIMEOUT; -import static org.apache.cassandra.config.DatabaseDescriptor.*; -import static org.apache.cassandra.db.ConsistencyLevel.*; import static org.apache.cassandra.locator.InetAddressAndPort.Serializer.inetAddressAndPortSerializer; import static org.apache.cassandra.locator.ReplicaLayout.forTokenWriteLiveAndDown; import static org.apache.cassandra.metrics.ClientRequestsMetricsHolder.casReadMetrics; @@ -126,9 +129,9 @@ import static org.apache.cassandra.service.paxos.Ballot.Flag.GLOBAL; import static org.apache.cassandra.service.paxos.Ballot.Flag.LOCAL; import static org.apache.cassandra.service.paxos.BallotGenerator.Global.nextBallot; import static org.apache.cassandra.service.paxos.BallotGenerator.Global.staleBallot; -import static org.apache.cassandra.service.paxos.ContentionStrategy.*; import static org.apache.cassandra.service.paxos.ContentionStrategy.Type.READ; import static org.apache.cassandra.service.paxos.ContentionStrategy.Type.WRITE; +import static org.apache.cassandra.service.paxos.ContentionStrategy.waitForContention; import static org.apache.cassandra.service.paxos.PaxosCommit.commit; import static org.apache.cassandra.service.paxos.PaxosCommitAndPrepare.commitAndPrepare; import static org.apache.cassandra.service.paxos.PaxosPrepare.prepare; @@ -439,7 +442,7 @@ public class Paxos } @Override - public void collectFailure(InetAddressAndPort inetAddressAndPort, RequestFailureReason t) + public void collectFailure(InetAddressAndPort inetAddressAndPort, RequestFailure t) { } diff --git a/src/java/org/apache/cassandra/service/paxos/PaxosCommit.java b/src/java/org/apache/cassandra/service/paxos/PaxosCommit.java index b5ce86794d..b79e032fe2 100644 --- a/src/java/org/apache/cassandra/service/paxos/PaxosCommit.java +++ b/src/java/org/apache/cassandra/service/paxos/PaxosCommit.java @@ -29,7 +29,7 @@ import org.apache.cassandra.config.CassandraRelevantProperties; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.ConsistencyLevel; import org.apache.cassandra.db.Mutation; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.locator.EndpointsForToken; import org.apache.cassandra.locator.InOurDc; import org.apache.cassandra.locator.InetAddressAndPort; @@ -44,13 +44,12 @@ import org.apache.cassandra.tracing.Tracing; import org.apache.cassandra.utils.concurrent.ConditionAsConsumer; import static java.util.Collections.emptyMap; -import static org.apache.cassandra.exceptions.RequestFailureReason.NODE_DOWN; import static org.apache.cassandra.exceptions.RequestFailureReason.UNKNOWN; import static org.apache.cassandra.net.Verb.PAXOS2_COMMIT_REMOTE_REQ; import static org.apache.cassandra.net.Verb.PAXOS_COMMIT_REQ; import static org.apache.cassandra.service.StorageProxy.shouldHint; import static org.apache.cassandra.service.StorageProxy.submitHint; -import static org.apache.cassandra.service.paxos.Commit.*; +import static org.apache.cassandra.service.paxos.Commit.Agreed; import static org.apache.cassandra.utils.concurrent.ConditionAsConsumer.newConditionAsConsumer; // Does not support EACH_QUORUM, as no such thing as EACH_SERIAL @@ -186,7 +185,7 @@ public class PaxosCommit> ex executeOnSelf |= isSelfOrSend(commitMessage, mutationMessage, participants.allLive.endpoint(i)); for (int i = 0, mi = participants.allDown.size(); i < mi ; ++i) - onFailure(participants.allDown.endpoint(i), NODE_DOWN); + onFailure(participants.allDown.endpoint(i), RequestFailure.NODE_DOWN); if (executeOnSelf) { @@ -223,7 +222,7 @@ public class PaxosCommit> ex * Record a failure or timeout, and maybe submit a hint to {@code from} */ @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason reason) + public void onFailure(InetAddressAndPort from, RequestFailure reason) { if (logger.isTraceEnabled()) logger.trace("{} {} from {}", commit, reason, from); diff --git a/src/java/org/apache/cassandra/service/paxos/PaxosPrepare.java b/src/java/org/apache/cassandra/service/paxos/PaxosPrepare.java index f293f1dba9..4a3d69eee1 100644 --- a/src/java/org/apache/cassandra/service/paxos/PaxosPrepare.java +++ b/src/java/org/apache/cassandra/service/paxos/PaxosPrepare.java @@ -34,9 +34,15 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.db.*; +import org.apache.cassandra.db.ColumnFamilyStore; +import org.apache.cassandra.db.DecoratedKey; +import org.apache.cassandra.db.Keyspace; +import org.apache.cassandra.db.ReadCommand; +import org.apache.cassandra.db.ReadExecutionController; +import org.apache.cassandra.db.ReadResponse; +import org.apache.cassandra.db.SinglePartitionReadCommand; import org.apache.cassandra.db.partitions.UnfilteredPartitionIterator; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.UnavailableException; import org.apache.cassandra.gms.EndpointState; import org.apache.cassandra.gms.Gossiper; @@ -44,7 +50,6 @@ import org.apache.cassandra.io.IVersionedSerializer; import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.io.util.DataOutputPlus; import org.apache.cassandra.locator.InetAddressAndPort; -import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.metrics.PaxosMetrics; import org.apache.cassandra.net.IVerbHandler; import org.apache.cassandra.net.Message; @@ -53,6 +58,7 @@ import org.apache.cassandra.schema.Schema; import org.apache.cassandra.schema.TableId; import org.apache.cassandra.schema.TableMetadata; import org.apache.cassandra.service.paxos.PaxosPrepare.Status.Outcome; +import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.tcm.ClusterMetadataService; import org.apache.cassandra.tcm.Epoch; import org.apache.cassandra.tracing.Tracing; @@ -64,13 +70,33 @@ import static org.apache.cassandra.locator.InetAddressAndPort.Serializer.inetAdd import static org.apache.cassandra.net.Verb.PAXOS2_PREPARE_REQ; import static org.apache.cassandra.net.Verb.PAXOS2_PREPARE_RSP; import static org.apache.cassandra.service.paxos.Ballot.Flag.NONE; -import static org.apache.cassandra.service.paxos.Commit.*; -import static org.apache.cassandra.service.paxos.Paxos.*; -import static org.apache.cassandra.service.paxos.PaxosPrepare.Status.Outcome.*; +import static org.apache.cassandra.service.paxos.Commit.Accepted; +import static org.apache.cassandra.service.paxos.Commit.Committed; +import static org.apache.cassandra.service.paxos.Commit.CompareResult; +import static org.apache.cassandra.service.paxos.Commit.isAfter; +import static org.apache.cassandra.service.paxos.Paxos.Electorate; +import static org.apache.cassandra.service.paxos.Paxos.LOG_TTL_LINEARIZABILITY_VIOLATIONS; +import static org.apache.cassandra.service.paxos.Paxos.Participants; +import static org.apache.cassandra.service.paxos.Paxos.consistency; +import static org.apache.cassandra.service.paxos.Paxos.getPaxosVariant; +import static org.apache.cassandra.service.paxos.Paxos.isInRangeAndShouldProcess; +import static org.apache.cassandra.service.paxos.Paxos.newBallot; +import static org.apache.cassandra.service.paxos.Paxos.verifyElectorate; +import static org.apache.cassandra.service.paxos.PaxosPrepare.Status.Outcome.ELECTORATE_MISMATCH; +import static org.apache.cassandra.service.paxos.PaxosPrepare.Status.Outcome.FOUND_INCOMPLETE_ACCEPTED; +import static org.apache.cassandra.service.paxos.PaxosPrepare.Status.Outcome.FOUND_INCOMPLETE_COMMITTED; +import static org.apache.cassandra.service.paxos.PaxosPrepare.Status.Outcome.MAYBE_FAILURE; +import static org.apache.cassandra.service.paxos.PaxosPrepare.Status.Outcome.PROMISED; +import static org.apache.cassandra.service.paxos.PaxosPrepare.Status.Outcome.READ_PERMITTED; +import static org.apache.cassandra.service.paxos.PaxosPrepare.Status.Outcome.SUPERSEDED; +import static org.apache.cassandra.service.paxos.PaxosState.MaybePromise; +import static org.apache.cassandra.service.paxos.PaxosState.MaybePromise.Outcome.PERMIT_READ; +import static org.apache.cassandra.service.paxos.PaxosState.MaybePromise.Outcome.PROMISE; +import static org.apache.cassandra.service.paxos.PaxosState.MaybePromise.Outcome.REJECT; +import static org.apache.cassandra.service.paxos.PaxosState.Snapshot; +import static org.apache.cassandra.service.paxos.PaxosState.get; import static org.apache.cassandra.utils.Clock.Global.currentTimeMillis; import static org.apache.cassandra.utils.Clock.Global.nanoTime; -import static org.apache.cassandra.service.paxos.PaxosState.*; -import static org.apache.cassandra.service.paxos.PaxosState.MaybePromise.Outcome.*; import static org.apache.cassandra.utils.CollectionSerializers.deserializeMap; import static org.apache.cassandra.utils.CollectionSerializers.serializeMap; import static org.apache.cassandra.utils.CollectionSerializers.serializedMapSize; @@ -804,7 +830,7 @@ public class PaxosPrepare extends PaxosRequestCallback im } @Override - public synchronized void onFailure(InetAddressAndPort from, RequestFailureReason reason) + public synchronized void onFailure(InetAddressAndPort from, RequestFailure reason) { if (logger.isTraceEnabled()) logger.trace("{} {} failure from {}", request, reason, from); @@ -812,7 +838,7 @@ public class PaxosPrepare extends PaxosRequestCallback im if (isDone()) return; - super.onFailureWithMutex(from, reason); + super.onFailureWithMutex(from, reason.reason); ++failures; if (failures + participants.sizeOfConsensusQuorum == 1 + participants.sizeOfPoll()) @@ -875,7 +901,7 @@ public class PaxosPrepare extends PaxosRequestCallback im } @Override - public void onRefreshFailure(InetAddressAndPort from, RequestFailureReason reason) + public void onRefreshFailure(InetAddressAndPort from, RequestFailure reason) { onFailure(from, reason); } diff --git a/src/java/org/apache/cassandra/service/paxos/PaxosPrepareRefresh.java b/src/java/org/apache/cassandra/service/paxos/PaxosPrepareRefresh.java index fbdab4c6fd..ad909c6a00 100644 --- a/src/java/org/apache/cassandra/service/paxos/PaxosPrepareRefresh.java +++ b/src/java/org/apache/cassandra/service/paxos/PaxosPrepareRefresh.java @@ -24,6 +24,7 @@ import java.util.List; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.exceptions.WriteTimeoutException; import org.apache.cassandra.io.IVersionedSerializer; @@ -38,8 +39,6 @@ import org.apache.cassandra.service.paxos.Commit.Agreed; import org.apache.cassandra.service.paxos.Commit.Committed; import org.apache.cassandra.tracing.Tracing; -import static org.apache.cassandra.exceptions.RequestFailureReason.TIMEOUT; -import static org.apache.cassandra.exceptions.RequestFailureReason.UNKNOWN; import static org.apache.cassandra.net.Verb.PAXOS2_PREPARE_REFRESH_REQ; import static org.apache.cassandra.service.paxos.Commit.isAfter; import static org.apache.cassandra.service.paxos.PaxosRequestCallback.shouldExecuteOnSelf; @@ -65,7 +64,7 @@ public class PaxosPrepareRefresh implements RequestCallbackWithFailure> } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason reason) + public void onFailure(InetAddressAndPort from, RequestFailure reason) { if (logger.isTraceEnabled()) logger.trace("{} {} failure from {}", proposal, reason, from); diff --git a/src/java/org/apache/cassandra/service/paxos/PaxosRepair.java b/src/java/org/apache/cassandra/service/paxos/PaxosRepair.java index 0242a5110e..ed369539ba 100644 --- a/src/java/org/apache/cassandra/service/paxos/PaxosRepair.java +++ b/src/java/org/apache/cassandra/service/paxos/PaxosRepair.java @@ -45,7 +45,7 @@ import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.db.Keyspace; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.UnavailableException; import org.apache.cassandra.io.IVersionedSerializer; import org.apache.cassandra.io.util.DataInputPlus; @@ -189,7 +189,7 @@ public class PaxosRepair extends AbstractPaxosRepair private Ballot clashingPromise; @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason reason) + public void onFailure(InetAddressAndPort from, RequestFailure reason) { updateState(this, null, (i1, i2) -> i1.onFailure()); } diff --git a/src/java/org/apache/cassandra/service/paxos/PaxosRequestCallback.java b/src/java/org/apache/cassandra/service/paxos/PaxosRequestCallback.java index aad32ace05..ff5f2d4068 100644 --- a/src/java/org/apache/cassandra/service/paxos/PaxosRequestCallback.java +++ b/src/java/org/apache/cassandra/service/paxos/PaxosRequestCallback.java @@ -24,14 +24,14 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.cassandra.config.CassandraRelevantProperties; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.WriteTimeoutException; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.net.Message; import org.apache.cassandra.service.FailureRecordingCallback; -import static org.apache.cassandra.exceptions.RequestFailureReason.TIMEOUT; -import static org.apache.cassandra.exceptions.RequestFailureReason.UNKNOWN; +import static org.apache.cassandra.exceptions.RequestFailure.TIMEOUT; +import static org.apache.cassandra.exceptions.RequestFailure.UNKNOWN; import static org.apache.cassandra.utils.FBUtilities.getBroadcastAddressAndPort; public abstract class PaxosRequestCallback extends FailureRecordingCallback @@ -58,7 +58,7 @@ public abstract class PaxosRequestCallback extends FailureRecordingCallback implements RequestCa } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason reason) + public void onFailure(InetAddressAndPort from, RequestFailure reason) { tryFailure(new PaxosCleanupException("Timed out waiting on response from " + from)); } diff --git a/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosCleanupSession.java b/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosCleanupSession.java index 80f571cd26..62b288986e 100644 --- a/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosCleanupSession.java +++ b/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosCleanupSession.java @@ -19,7 +19,9 @@ package org.apache.cassandra.service.paxos.cleanup; import java.lang.ref.WeakReference; -import java.util.*; +import java.util.Collection; +import java.util.Queue; +import java.util.UUID; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; @@ -28,7 +30,7 @@ import com.google.common.base.Preconditions; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.gms.ApplicationState; import org.apache.cassandra.gms.EndpointState; import org.apache.cassandra.gms.IEndpointStateChangeSubscriber; @@ -239,7 +241,7 @@ public class PaxosCleanupSession extends AsyncFuture implements Runnable, } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason reason) + public void onFailure(InetAddressAndPort from, RequestFailure reason) { fail(from.toString() + ' ' + reason + " for cleanup request for paxos cleanup session " + session); } diff --git a/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosFinishPrepareCleanup.java b/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosFinishPrepareCleanup.java index 07b1bbe334..0fc189f2e4 100644 --- a/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosFinishPrepareCleanup.java +++ b/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosFinishPrepareCleanup.java @@ -18,9 +18,11 @@ package org.apache.cassandra.service.paxos.cleanup; -import java.util.*; +import java.util.Collection; +import java.util.HashSet; +import java.util.Set; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.net.IVerbHandler; import org.apache.cassandra.net.Message; @@ -51,7 +53,7 @@ public class PaxosFinishPrepareCleanup extends AsyncFuture implements Requ } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason reason) + public void onFailure(InetAddressAndPort from, RequestFailure reason) { tryFailure(new PaxosCleanupException(reason + " failure response from " + from)); } diff --git a/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosStartPrepareCleanup.java b/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosStartPrepareCleanup.java index a375481c34..08cbd9aa09 100644 --- a/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosStartPrepareCleanup.java +++ b/src/java/org/apache/cassandra/service/paxos/cleanup/PaxosStartPrepareCleanup.java @@ -19,17 +19,23 @@ package org.apache.cassandra.service.paxos.cleanup; import java.io.IOException; -import java.util.*; +import java.util.ArrayList; +import java.util.Collection; +import java.util.HashSet; +import java.util.List; +import java.util.Set; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.apache.cassandra.db.*; +import org.apache.cassandra.db.ColumnFamilyStore; +import org.apache.cassandra.db.ConsistencyLevel; +import org.apache.cassandra.db.TypeSizes; import org.apache.cassandra.dht.AbstractBounds; import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.gms.EndpointState; import org.apache.cassandra.gms.HeartBeatState; import org.apache.cassandra.io.IVersionedSerializer; @@ -94,7 +100,7 @@ public class PaxosStartPrepareCleanup extends AsyncFuture i } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason reason) + public void onFailure(InetAddressAndPort from, RequestFailure reason) { tryFailure(new PaxosCleanupException("Received " + reason + " failure response from " + from)); } diff --git a/src/java/org/apache/cassandra/service/reads/ReadCallback.java b/src/java/org/apache/cassandra/service/reads/ReadCallback.java index ca25d1a3fd..e266de0616 100644 --- a/src/java/org/apache/cassandra/service/reads/ReadCallback.java +++ b/src/java/org/apache/cassandra/service/reads/ReadCallback.java @@ -34,6 +34,7 @@ import org.apache.cassandra.db.ReadCommand; import org.apache.cassandra.db.ReadResponse; import org.apache.cassandra.exceptions.ReadFailureException; import org.apache.cassandra.exceptions.ReadTimeoutException; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.locator.Endpoints; import org.apache.cassandra.locator.InetAddressAndPort; @@ -186,7 +187,7 @@ public class ReadCallback, P extends ReplicaPlan.ForRead< InetAddressAndPort from = message.from(); if (WarningContext.isSupported(params.keySet())) { - RequestFailureReason reason = getWarningContext().updateCounters(params, from); + RequestFailure reason = getWarningContext().updateCounters(params, from); replicaPlan().collectFailure(message.from(), reason); if (reason != null) { @@ -236,11 +237,11 @@ public class ReadCallback, P extends ReplicaPlan.ForRead< } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { assertWaitingFor(from); - failureReasonByEndpoint.put(from, failureReason); + failureReasonByEndpoint.put(from, failure.reason); if (replicaPlan().readQuorum() + failuresUpdater.incrementAndGet(this) > replicaPlan().contacts().size()) condition.signalAll(); diff --git a/src/java/org/apache/cassandra/service/reads/thresholds/WarningContext.java b/src/java/org/apache/cassandra/service/reads/thresholds/WarningContext.java index dd6ee2f1a6..5bb5deb99a 100644 --- a/src/java/org/apache/cassandra/service/reads/thresholds/WarningContext.java +++ b/src/java/org/apache/cassandra/service/reads/thresholds/WarningContext.java @@ -22,7 +22,7 @@ import java.util.EnumSet; import java.util.Map; import java.util.Set; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.net.ParamType; @@ -43,31 +43,31 @@ public class WarningContext return !Collections.disjoint(keys, SUPPORTED); } - public RequestFailureReason updateCounters(Map params, InetAddressAndPort from) + public RequestFailure updateCounters(Map params, InetAddressAndPort from) { for (Map.Entry entry : params.entrySet()) { WarnAbortCounter counter = null; - RequestFailureReason reason = null; + RequestFailure reason = null; switch (entry.getKey()) { case ROW_INDEX_READ_SIZE_FAIL: - reason = RequestFailureReason.READ_SIZE; + reason = RequestFailure.READ_SIZE; case ROW_INDEX_READ_SIZE_WARN: counter = rowIndexReadSize; break; case LOCAL_READ_SIZE_FAIL: - reason = RequestFailureReason.READ_SIZE; + reason = RequestFailure.READ_SIZE; case LOCAL_READ_SIZE_WARN: counter = localReadSize; break; case TOMBSTONE_FAIL: - reason = RequestFailureReason.READ_TOO_MANY_TOMBSTONES; + reason = RequestFailure.READ_TOO_MANY_TOMBSTONES; case TOMBSTONE_WARNING: counter = tombstones; break; case TOO_MANY_REFERENCED_INDEXES_FAIL: - reason = RequestFailureReason.READ_TOO_MANY_INDEXES; + reason = RequestFailure.READ_TOO_MANY_INDEXES; case TOO_MANY_REFERENCED_INDEXES_WARN: counter = indexReadSSTablesCount; break; diff --git a/src/java/org/apache/cassandra/tcm/PaxosBackedProcessor.java b/src/java/org/apache/cassandra/tcm/PaxosBackedProcessor.java index 19bdc7ae9b..45b5945cbc 100644 --- a/src/java/org/apache/cassandra/tcm/PaxosBackedProcessor.java +++ b/src/java/org/apache/cassandra/tcm/PaxosBackedProcessor.java @@ -32,8 +32,8 @@ import org.slf4j.LoggerFactory; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.ConsistencyLevel; import org.apache.cassandra.exceptions.ReadTimeoutException; -import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.locator.EndpointsForRange; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.locator.Replica; import org.apache.cassandra.metrics.TCMMetrics; @@ -202,10 +202,10 @@ public class PaxosBackedProcessor extends AbstractLocalProcessor } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { - logger.debug("Error response from {} with {}", from, failureReason); - condition.tryFailure(new TimeoutException(failureReason.toString())); + logger.debug("Error response from {} with {}", from, failure.reason); + condition.tryFailure(new TimeoutException(failure.reason.toString())); } public void retry() diff --git a/src/java/org/apache/cassandra/tcm/RemoteProcessor.java b/src/java/org/apache/cassandra/tcm/RemoteProcessor.java index 0ea055b908..e5cb0568fe 100644 --- a/src/java/org/apache/cassandra/tcm/RemoteProcessor.java +++ b/src/java/org/apache/cassandra/tcm/RemoteProcessor.java @@ -34,6 +34,7 @@ import org.slf4j.LoggerFactory; import com.codahale.metrics.Timer; import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.gms.FailureDetector; import org.apache.cassandra.locator.InetAddressAndPort; @@ -213,7 +214,7 @@ public final class RemoteProcessor implements Processor (attempt, from, failure) -> { if (promise.isDone() || promise.isCancelled()) return false; - if (failure == RequestFailureReason.NOT_CMS) + if (failure.reason == RequestFailureReason.NOT_CMS) { logger.debug("{} is not a member of the CMS, querying it to discover current membership", from); DiscoveredNodes cms = tryDiscover(from); @@ -257,7 +258,7 @@ public final class RemoteProcessor implements Processor } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { // "success" - this lets us just try the next one in cmsIter promise.setSuccess(new DiscoveredNodes(Collections.emptySet(), DiscoveredNodes.Kind.KNOWN_PEERS)); diff --git a/src/java/org/apache/cassandra/tcm/migration/Election.java b/src/java/org/apache/cassandra/tcm/migration/Election.java index 507a55d31c..6ada116323 100644 --- a/src/java/org/apache/cassandra/tcm/migration/Election.java +++ b/src/java/org/apache/cassandra/tcm/migration/Election.java @@ -29,7 +29,6 @@ import java.util.concurrent.atomic.AtomicReference; import java.util.stream.Collectors; import com.google.common.collect.Sets; - import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/src/java/org/apache/cassandra/tcm/sequences/ProgressBarrier.java b/src/java/org/apache/cassandra/tcm/sequences/ProgressBarrier.java index af504d35d3..8212a57d9d 100644 --- a/src/java/org/apache/cassandra/tcm/sequences/ProgressBarrier.java +++ b/src/java/org/apache/cassandra/tcm/sequences/ProgressBarrier.java @@ -39,7 +39,7 @@ import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.ConsistencyLevel; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.locator.EndpointsForRange; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.locator.Replica; @@ -544,10 +544,10 @@ public class ProgressBarrier } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { - logger.debug("Error response from {} with {}", from, failureReason); - condition.tryFailure(new TimeoutException(String.format("Watermark request did returned %s.", failureReason))); + logger.debug("Error response from {} with {}", from, failure); + condition.tryFailure(new TimeoutException(String.format("Watermark request did returned %s.", failure.reason))); } public void retry() diff --git a/src/java/org/apache/cassandra/utils/FBUtilities.java b/src/java/org/apache/cassandra/utils/FBUtilities.java index 6892777422..7417f757bc 100644 --- a/src/java/org/apache/cassandra/utils/FBUtilities.java +++ b/src/java/org/apache/cassandra/utils/FBUtilities.java @@ -94,7 +94,6 @@ import org.apache.cassandra.utils.concurrent.FutureCombiner; import org.apache.cassandra.utils.concurrent.UncheckedInterruptedException; import org.objectweb.asm.Opcodes; -import static org.apache.cassandra.config.CassandraRelevantProperties.CASSANDRA_AVAILABLE_PROCESSORS; import static org.apache.cassandra.config.CassandraRelevantProperties.BUILD_DATE; import static org.apache.cassandra.config.CassandraRelevantProperties.GIT_SHA; import static org.apache.cassandra.config.CassandraRelevantProperties.LINE_SEPARATOR; @@ -131,8 +130,6 @@ public class FBUtilities private static volatile String previousReleaseVersionString; - private static final int availableProcessors = CASSANDRA_AVAILABLE_PROCESSORS.getInt(DatabaseDescriptor.getAvailableProcessors()); - private static volatile Supplier systemInfoSupplier = Suppliers.memoize(SystemInfo::new); public static void setAvailableProcessors(int value) diff --git a/test/data/config/version=5.0-alpha1.yml b/test/data/config/version=5.0-alpha1.yml index 8dad0f60ac..19995ce52b 100644 --- a/test/data/config/version=5.0-alpha1.yml +++ b/test/data/config/version=5.0-alpha1.yml @@ -407,7 +407,7 @@ max_concurrent_automatic_sstable_upgrades: "java.lang.Integer" maximum_replication_factor_warn_threshold: "java.lang.Integer" denylist_reads_enabled: "java.lang.Boolean" permissions_cache_active_update: "java.lang.Boolean" -available_processors: "java.lang.Integer" +available_processors: "org.apache.cassandra.config.OptionaldPositiveInt" file_cache_round_up: "java.lang.Boolean" secondary_indexes_per_table_warn_threshold: "java.lang.Integer" tables_warn_threshold: "java.lang.Integer" diff --git a/test/simulator/main/org/apache/cassandra/simulator/systems/SimulatedAction.java b/test/simulator/main/org/apache/cassandra/simulator/systems/SimulatedAction.java index ca5407fe89..b30626136b 100644 --- a/test/simulator/main/org/apache/cassandra/simulator/systems/SimulatedAction.java +++ b/test/simulator/main/org/apache/cassandra/simulator/systems/SimulatedAction.java @@ -32,7 +32,7 @@ import com.google.common.base.Preconditions; import org.apache.cassandra.concurrent.ImmediateExecutor; import org.apache.cassandra.distributed.api.IInvokableInstance; import org.apache.cassandra.distributed.api.IMessage; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.net.RequestCallback; import org.apache.cassandra.net.RequestCallbacks; @@ -396,7 +396,7 @@ public abstract class SimulatedAction extends Action implements InterceptorOfCon if (callback != null) { RequestCallback invokeOn = (RequestCallback) callback.callback; - RequestFailureReason reason = innerIsTimeout ? RequestFailureReason.TIMEOUT : RequestFailureReason.UNKNOWN; + RequestFailure reason = innerIsTimeout ? RequestFailure.TIMEOUT : RequestFailure.UNKNOWN; invokeOn.onFailure(address, reason); } return null; diff --git a/test/unit/org/apache/cassandra/config/ConfigCompatibilityTest.java b/test/unit/org/apache/cassandra/config/ConfigCompatibilityTest.java index f5d10a6b52..2d38db6530 100644 --- a/test/unit/org/apache/cassandra/config/ConfigCompatibilityTest.java +++ b/test/unit/org/apache/cassandra/config/ConfigCompatibilityTest.java @@ -119,6 +119,7 @@ public class ConfigCompatibilityTest .add("Property role_manager used to be a value-type, but now is nested type class org.apache.cassandra.config.ParameterizedClass") .add("Property network_authorizer used to be a value-type, but now is nested type class org.apache.cassandra.config.ParameterizedClass") .add("require_client_auth types do not match; java.lang.String != java.lang.Boolean") + .add("available_processors types do not match; org.apache.cassandra.config.OptionaldPositiveInt != java.lang.Integer") .build(); /** diff --git a/test/unit/org/apache/cassandra/db/CounterMutationVerbHandlerOutOfRangeTest.java b/test/unit/org/apache/cassandra/db/CounterMutationVerbHandlerOutOfRangeTest.java index 9ae05e3e4e..da58eccfa0 100644 --- a/test/unit/org/apache/cassandra/db/CounterMutationVerbHandlerOutOfRangeTest.java +++ b/test/unit/org/apache/cassandra/db/CounterMutationVerbHandlerOutOfRangeTest.java @@ -37,7 +37,7 @@ import org.apache.cassandra.db.rows.Cell; import org.apache.cassandra.db.rows.Row; import org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper; import org.apache.cassandra.exceptions.InvalidRoutingException; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.metrics.StorageMetrics; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessagingService; @@ -51,12 +51,18 @@ import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.tcm.membership.NodeState; import org.apache.cassandra.utils.FBUtilities; -import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.*; -import static org.apache.cassandra.utils.ByteBufferUtil.bytes; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.MessageDelivery; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.broadcastAddress; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.bytesToken; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.node1; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.randomInt; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.registerOutgoingMessageSink; +import static org.apache.cassandra.utils.ByteBufferUtil.bytes; + public class CounterMutationVerbHandlerOutOfRangeTest { private static final String KEYSPACE = "CounterCacheTest"; @@ -168,7 +174,7 @@ public class CounterMutationVerbHandlerOutOfRangeTest MessageDelivery response = messageSink.get(100, TimeUnit.MILLISECONDS); assertEquals(Verb.FAILURE_RSP, response.message.verb()); assertEquals(broadcastAddress, response.message.from()); - assertTrue(response.message.payload instanceof RequestFailureReason); + assertTrue(response.message.payload instanceof RequestFailure); assertEquals(messageId, response.message.id()); assertEquals(node1, response.to); } diff --git a/test/unit/org/apache/cassandra/db/MutationVerbHandlerOutOfRangeTest.java b/test/unit/org/apache/cassandra/db/MutationVerbHandlerOutOfRangeTest.java index 0421e26aa5..77ef47192e 100644 --- a/test/unit/org/apache/cassandra/db/MutationVerbHandlerOutOfRangeTest.java +++ b/test/unit/org/apache/cassandra/db/MutationVerbHandlerOutOfRangeTest.java @@ -37,7 +37,7 @@ import org.apache.cassandra.db.rows.Cell; import org.apache.cassandra.db.rows.Row; import org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper; import org.apache.cassandra.exceptions.InvalidRoutingException; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.metrics.StorageMetrics; import org.apache.cassandra.net.IVerbHandler; import org.apache.cassandra.net.Message; @@ -51,11 +51,17 @@ import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.tcm.membership.NodeState; import org.apache.cassandra.utils.FBUtilities; -import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.*; -import static org.apache.cassandra.utils.ByteBufferUtil.bytes; import static org.junit.Assert.assertEquals; import static org.junit.Assert.fail; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.MessageDelivery; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.broadcastAddress; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.bytesToken; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.node1; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.randomInt; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.registerOutgoingMessageSink; +import static org.apache.cassandra.utils.ByteBufferUtil.bytes; + public class MutationVerbHandlerOutOfRangeTest { private static final String TEST_NAME = "mutation_vh_test_"; @@ -190,7 +196,7 @@ public class MutationVerbHandlerOutOfRangeTest MessageDelivery response = messageSink.get(100, TimeUnit.MILLISECONDS); assertEquals(isOutOfRange ? Verb.FAILURE_RSP : Verb.MUTATION_RSP, response.message.verb()); assertEquals(broadcastAddress, response.message.from()); - assertEquals(isOutOfRange, response.message.payload instanceof RequestFailureReason); + assertEquals(isOutOfRange, response.message.payload instanceof RequestFailure); assertEquals(messageId, response.message.id()); assertEquals(node1, response.to); assertEquals(startingTotalMetricCount + (isOutOfRange ? 1 : 0), StorageMetrics.totalOpsForInvalidToken.getCount()); diff --git a/test/unit/org/apache/cassandra/db/ReadCommandVerbHandlerOutOfRangeTest.java b/test/unit/org/apache/cassandra/db/ReadCommandVerbHandlerOutOfRangeTest.java index 847326ca4d..419772d304 100644 --- a/test/unit/org/apache/cassandra/db/ReadCommandVerbHandlerOutOfRangeTest.java +++ b/test/unit/org/apache/cassandra/db/ReadCommandVerbHandlerOutOfRangeTest.java @@ -37,7 +37,7 @@ import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; import org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper; import org.apache.cassandra.exceptions.InvalidRoutingException; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.metrics.StorageMetrics; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessagingService; @@ -50,10 +50,16 @@ import org.apache.cassandra.tcm.membership.NodeState; import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.FBUtilities; -import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.*; -import static org.apache.cassandra.net.Verb.READ_REQ; import static org.junit.Assert.assertEquals; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.MessageDelivery; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.broadcastAddress; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.bytesToken; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.node1; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.randomInt; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.registerOutgoingMessageSink; +import static org.apache.cassandra.net.Verb.READ_REQ; + public class ReadCommandVerbHandlerOutOfRangeTest { private static ReadCommandVerbHandler handler; @@ -184,7 +190,7 @@ public class ReadCommandVerbHandlerOutOfRangeTest MessageDelivery response = messageSink.get(100, TimeUnit.MILLISECONDS); assertEquals(isOutOfRange ? Verb.FAILURE_RSP : Verb.READ_RSP, response.message.verb()); assertEquals(broadcastAddress, response.message.from()); - assertEquals(isOutOfRange, response.message.payload instanceof RequestFailureReason); + assertEquals(isOutOfRange, response.message.payload instanceof RequestFailure); assertEquals(messageId, response.message.id()); assertEquals(node1, response.to); assertEquals(startingTotalMetricCount + (isOutOfRange ? 1 : 0), StorageMetrics.totalOpsForInvalidToken.getCount()); diff --git a/test/unit/org/apache/cassandra/exceptions/RemoteExceptionTest.java b/test/unit/org/apache/cassandra/exceptions/RemoteExceptionTest.java new file mode 100644 index 0000000000..42e211e162 --- /dev/null +++ b/test/unit/org/apache/cassandra/exceptions/RemoteExceptionTest.java @@ -0,0 +1,112 @@ +/* + * 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.exceptions; + +import java.util.HashMap; +import java.util.Map; + +import org.junit.Test; + +import org.apache.cassandra.io.util.DataInputBuffer; +import org.apache.cassandra.io.util.DataOutputBuffer; +import org.apache.cassandra.net.MessagingService; + +import static com.google.common.base.Throwables.getStackTraceAsString; +import static org.apache.cassandra.exceptions.ExceptionSerializer.getMessageWithOriginatingHost; +import static org.apache.cassandra.exceptions.ExceptionSerializer.nullableRemoteExceptionSerializer; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNull; + +public class RemoteExceptionTest +{ + @Test + public void testRoundtrip() throws Exception + { + testRoundtrip(null); + Throwable root = new Throwable(); + testRoundtrip(root); + Throwable suppressed = new Throwable(); + Throwable causedByRoot = new Throwable(root); + testRoundtrip(causedByRoot); + causedByRoot.addSuppressed(root); + testRoundtrip(causedByRoot); + root.addSuppressed(causedByRoot); + testRoundtrip(root); + root.addSuppressed(suppressed); + testRoundtrip(root); + } + + public void testRoundtrip(Throwable original) throws Exception + { + Throwable normalizedOriginal = normalizeThrowable(original); + + DataOutputBuffer dob = new DataOutputBuffer(); + nullableRemoteExceptionSerializer.serialize(original, dob, MessagingService.current_version); + assertEquals(nullableRemoteExceptionSerializer.serializedSize(original, MessagingService.current_version), dob.toByteArray().length); + DataInputBuffer dib = new DataInputBuffer(dob.toByteArray()); + Throwable test = nullableRemoteExceptionSerializer.deserialize(dib, MessagingService.current_version); + if (original == null) + { + assertNull(test); + } + else + { + String originalString = getStackTraceAsString(normalizedOriginal); + String testString = getStackTraceAsString(test); + assertEquals(originalString, testString); + } + } + + public static Throwable normalizeThrowable(Throwable t) throws Exception + { + return normalizeThrowable(t, true, new HashMap<>()); + } + + private static Throwable normalizeThrowable(Throwable t, boolean isFirstException, Map alreadyNormalized) throws Exception + { + if (t == null) + return null; + + if (alreadyNormalized.containsKey(t)) + return alreadyNormalized.get(t); + + // Classloader, module name, and module version are difficult to get right because STE doesn't + // expose enough parameters to serialize the formatting correctly so settle for something close, but not exact + // Alternatives look fragile across different JVM versions and yield only moderate additional debugability + // when using class loaders and modules + StackTraceElement[] originalStack = t.getStackTrace(); + StackTraceElement[] normalizedStack = new StackTraceElement[originalStack.length]; + for (int i = 0; i < originalStack.length; i++) + { + StackTraceElement originalSTE = originalStack[i]; + normalizedStack[i] = new StackTraceElement(originalSTE.getClassName(), originalSTE.getMethodName(), originalSTE.getFileName(), originalSTE.getLineNumber()); + } + + Throwable normalized; + if (t.getCause() == null) + normalized = t.getClass().getConstructor(String.class).newInstance(getMessageWithOriginatingHost(t, isFirstException)); + else + normalized = t.getClass().getConstructor(String.class, Throwable.class).newInstance(getMessageWithOriginatingHost(t, isFirstException), normalizeThrowable(t.getCause(), false, alreadyNormalized)); + alreadyNormalized.put(t, normalized); + normalized.setStackTrace(normalizedStack); + for (Throwable suppressed : t.getSuppressed()) + normalized.addSuppressed(normalizeThrowable(suppressed, false, alreadyNormalized)); + return normalized; + } +} \ No newline at end of file diff --git a/test/unit/org/apache/cassandra/net/ConnectionTest.java b/test/unit/org/apache/cassandra/net/ConnectionTest.java index 42c137cc5f..70bb0c8046 100644 --- a/test/unit/org/apache/cassandra/net/ConnectionTest.java +++ b/test/unit/org/apache/cassandra/net/ConnectionTest.java @@ -59,7 +59,7 @@ import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.config.EncryptionOptions; import org.apache.cassandra.db.commitlog.CommitLog; import org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.UnknownColumnException; import org.apache.cassandra.io.IVersionedAsymmetricSerializer; import org.apache.cassandra.io.IVersionedSerializer; @@ -78,7 +78,7 @@ import static org.apache.cassandra.net.NoPayload.noPayload; import static org.apache.cassandra.net.MessagingService.current_version; import static org.apache.cassandra.net.ConnectionType.LARGE_MESSAGES; import static org.apache.cassandra.net.ConnectionType.SMALL_MESSAGES; -import static org.apache.cassandra.net.ConnectionUtils.*; +import static org.apache.cassandra.net.ConnectionUtils.check; import static org.apache.cassandra.net.OutboundConnectionSettings.Framing.LZ4; import static org.apache.cassandra.net.OutboundConnections.LARGE_MESSAGE_THRESHOLD; import static org.apache.cassandra.utils.Clock.Global.nanoTime; @@ -388,7 +388,7 @@ public class ConnectionTest MessagingService.instance().callbacks.addWithExpiration(new RequestCallback() { @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { done.countDown(); } diff --git a/test/unit/org/apache/cassandra/net/MessageDeliveryTest.java b/test/unit/org/apache/cassandra/net/MessageDeliveryTest.java index c2d4656ddc..e8fcf286a6 100644 --- a/test/unit/org/apache/cassandra/net/MessageDeliveryTest.java +++ b/test/unit/org/apache/cassandra/net/MessageDeliveryTest.java @@ -35,7 +35,7 @@ import org.apache.cassandra.concurrent.ScheduledExecutorPlus; import org.apache.cassandra.concurrent.SimulatedExecutorFactory; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.dht.Murmur3Partitioner; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.net.MessageDelivery.FailedResponseException; import org.apache.cassandra.net.MessageDelivery.MaxRetriesException; @@ -170,7 +170,7 @@ public class MessageDeliveryTest assertThat(result).isDone(); FailedResponseException e = getFailedResponseException(result); assertThat(e.from).isEqualTo(ID1); - assertThat(e.failure).isEqualTo(RequestFailureReason.TIMEOUT); + assertThat(e.failure).isEqualTo(RequestFailure.TIMEOUT); Mockito.verify(backoff, Mockito.times(1)).mayRetry(Mockito.anyInt()); Mockito.verify(backoff, Mockito.never()).computeWaitTime(Mockito.anyInt()); Mockito.verify(backoff, Mockito.never()).unit(); diff --git a/test/unit/org/apache/cassandra/net/MessageTest.java b/test/unit/org/apache/cassandra/net/MessageTest.java index 8e89973aa7..ddc5f6b9c6 100644 --- a/test/unit/org/apache/cassandra/net/MessageTest.java +++ b/test/unit/org/apache/cassandra/net/MessageTest.java @@ -19,8 +19,8 @@ package org.apache.cassandra.net; import java.io.IOException; import java.nio.ByteBuffer; -import java.nio.charset.CharacterCodingException; import java.nio.charset.StandardCharsets; +import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import java.util.stream.Stream; @@ -30,6 +30,7 @@ import org.junit.Test; import org.apache.cassandra.ServerTestUtils; import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.io.IVersionedSerializer; import org.apache.cassandra.io.util.DataInputBuffer; @@ -37,6 +38,7 @@ import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.io.util.DataOutputBuffer; import org.apache.cassandra.io.util.DataOutputPlus; import org.apache.cassandra.locator.InetAddressAndPort; +import org.apache.cassandra.net.MessagingService.Version; import org.apache.cassandra.tcm.Epoch; import org.apache.cassandra.tracing.Tracing; import org.apache.cassandra.tracing.Tracing.TraceType; @@ -44,16 +46,22 @@ import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.FreeRunningClock; import org.apache.cassandra.utils.TimeUUID; +import static com.google.common.base.Throwables.getStackTraceAsString; +import static org.apache.cassandra.exceptions.RemoteExceptionTest.normalizeThrowable; import static org.apache.cassandra.net.Message.serializer; import static org.apache.cassandra.net.MessagingService.VERSION_40; +import static org.apache.cassandra.net.MessagingService.VERSION_50; import static org.apache.cassandra.net.NoPayload.noPayload; import static org.apache.cassandra.net.ParamType.RESPOND_TO; import static org.apache.cassandra.net.ParamType.TRACE_SESSION; import static org.apache.cassandra.net.ParamType.TRACE_TYPE; import static org.apache.cassandra.utils.MonotonicClock.Global.approxTime; import static org.apache.cassandra.utils.TimeUUID.Generator.nextTimeUUID; - -import static org.junit.Assert.*; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertSame; +import static org.junit.Assert.assertTrue; public class MessageTest { @@ -165,7 +173,7 @@ public class MessageTest } @Test - public void testCycleNoPayload() throws IOException + public void testCycleNoPayload() throws Exception { Message msg = Message.builder(Verb._TEST_1, noPayload) @@ -190,15 +198,20 @@ public class MessageTest } @Test - public void testFailureResponse() throws IOException + public void testFailureResponse() throws Exception { long expiresAt = approxTime.now(); - Message msg = Message.failureResponse(1, expiresAt, RequestFailureReason.INCOMPATIBLE_SCHEMA); + ExecutionException cause = new ExecutionException("test", new NullPointerException()); + Throwable root = new Throwable(cause); + Throwable suppressed = new Throwable(); + root.addSuppressed(suppressed); + Message msg = Message.failureResponse(1, expiresAt, new RequestFailure(RequestFailureReason.INCOMPATIBLE_SCHEMA, root)); assertEquals(1, msg.id()); assertEquals(Verb.FAILURE_RSP, msg.verb()); assertEquals(expiresAt, msg.expiresAtNanos()); - assertEquals(RequestFailureReason.INCOMPATIBLE_SCHEMA, msg.payload); + assertEquals(RequestFailureReason.INCOMPATIBLE_SCHEMA, msg.payload.reason); + assertEquals(getStackTraceAsString(root), getStackTraceAsString(msg.payload.failure)); assertTrue(msg.isFailureResponse()); testCycle(msg); @@ -218,19 +231,26 @@ public class MessageTest } @Test - public void testCustomParams() throws CharacterCodingException, IOException + public void testCustomParams() throws IOException + { + for (Version version : MessagingService.Version.values()) + if (version.value >= VERSION_40) + testCustomParams(version.value); + } + + private void testCustomParams(int version) throws IOException { long id = 1; InetAddressAndPort from = FBUtilities.getLocalAddressAndPort(); Message msg = - Message.builder(Verb._TEST_1, noPayload) - .withEpoch(Epoch.EMPTY) - .withId(1) - .from(from) - .withCustomParam("custom1", "custom1value".getBytes(StandardCharsets.UTF_8)) - .withCustomParam("custom2", "custom2value".getBytes(StandardCharsets.UTF_8)) - .build(); + Message.builder(Verb._TEST_1, noPayload) + .withEpoch(Epoch.EMPTY) + .withId(1) + .from(from) + .withCustomParam("custom1", "custom1value".getBytes(StandardCharsets.UTF_8)) + .withCustomParam("custom2", "custom2value".getBytes(StandardCharsets.UTF_8)) + .build(); assertEquals(id, msg.id()); assertEquals(from, msg.from()); @@ -239,9 +259,10 @@ public class MessageTest assertEquals("custom2value", new String(msg.header.customParams().get("custom2"), StandardCharsets.UTF_8)); DataOutputBuffer out = DataOutputBuffer.scratchBuffer.get(); - Message.serializer.serialize(msg, out, VERSION_40); + out.clear(); + Message.serializer.serialize(msg, out, version); DataInputBuffer in = new DataInputBuffer(out.buffer(), true); - msg = Message.serializer.deserialize(in, from, VERSION_40); + msg = Message.serializer.deserialize(in, from, version); assertEquals(id, msg.id()); assertEquals(from, msg.from()); @@ -265,13 +286,13 @@ public class MessageTest } } - private void testCycle(Message msg) throws IOException + private void testCycle(Message msg) throws Exception { testCycle(msg, VERSION_40); } // serialize (using both variants, all in one or header then rest), verify serialized size, deserialize, compare to the original - private void testCycle(Message msg, int version) throws IOException + private void testCycle(Message msg, int version) throws Exception { try (DataOutputBuffer out = new DataOutputBuffer()) { @@ -283,7 +304,7 @@ public class MessageTest { Message msgOut = serializer.deserialize(in, msg.from(), version); assertEquals(0, in.available()); - assertMessagesEqual(msg, msgOut); + assertMessagesEqual(msg, msgOut, version); } // extract header first, then deserialize the rest of the message and compare outcomes @@ -293,12 +314,12 @@ public class MessageTest Message.Header headerOut = serializer.extractHeader(buffer, msg.from(), approxTime.now(), version); Message msgOut = serializer.deserialize(in, headerOut, version); assertEquals(0, in.available()); - assertMessagesEqual(msg, msgOut); + assertMessagesEqual(msg, msgOut, version); } } } - private static void assertMessagesEqual(Message msg1, Message msg2) + private static void assertMessagesEqual(Message msg1, Message msg2, int version) throws Exception { assertEquals(msg1.id(), msg2.id()); assertEquals(msg1.verb(), msg2.verb()); @@ -316,6 +337,19 @@ public class MessageTest assertTrue(payload2 == noPayload || payload2 == null); else if (null == payload2) assertSame(payload1, noPayload); + else if (msg1.verb() == Verb.FAILURE_RSP) + { + RequestFailure reason1 = (RequestFailure)msg1.payload; + RequestFailure reason2 = (RequestFailure)msg2.payload; + assertEquals(reason1.reason, reason2.reason); + if (version >= VERSION_50) + { + if (reason1.failure == null) + assertNull(reason2.failure); + else + assertEquals(getStackTraceAsString(normalizeThrowable(reason1.failure)), getStackTraceAsString(reason2.failure)); + } + } else assertEquals(payload1, payload2); } diff --git a/test/unit/org/apache/cassandra/net/SimulatedMessageDelivery.java b/test/unit/org/apache/cassandra/net/SimulatedMessageDelivery.java index 6a335f4aab..6cbabe37cc 100644 --- a/test/unit/org/apache/cassandra/net/SimulatedMessageDelivery.java +++ b/test/unit/org/apache/cassandra/net/SimulatedMessageDelivery.java @@ -31,7 +31,7 @@ import javax.annotation.Nullable; import accord.utilsfork.Gens; import accord.utilsfork.RandomSource; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.utils.concurrent.AsyncPromise; import org.apache.cassandra.utils.concurrent.Future; @@ -184,7 +184,7 @@ public class SimulatedMessageDelivery implements MessageDelivery } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failure) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { promise.tryFailure(new MessagingService.FailureResponseException(from, failure)); } @@ -237,7 +237,7 @@ public class SimulatedMessageDelivery implements MessageDelivery if (action == Action.FAILURE) onDropped.onDrop(action, to, message); if (callback != null) - scheduler.schedule(() -> callback.onFailure(to, RequestFailureReason.UNKNOWN), + scheduler.schedule(() -> callback.onFailure(to, RequestFailure.UNKNOWN), message.verb().expiresAfterNanos(), TimeUnit.NANOSECONDS); return; default: @@ -252,7 +252,7 @@ public class SimulatedMessageDelivery implements MessageDelivery assert ctx == cb; try { - ctx.onFailure(to, RequestFailureReason.TIMEOUT); + ctx.onFailure(to, RequestFailure.TIMEOUT); } catch (Throwable t) { @@ -302,7 +302,7 @@ public class SimulatedMessageDelivery implements MessageDelivery try { if (msg.isFailureResponse()) - callback.onFailure(msg.from(), (RequestFailureReason) msg.payload); + callback.onFailure(msg.from(), (RequestFailure) msg.payload); else callback.onResponse(msg); } catch (Throwable t) @@ -364,7 +364,7 @@ public class SimulatedMessageDelivery implements MessageDelivery callback.onResponse(msg); } - public void onFailure(InetAddressAndPort from, RequestFailureReason failure) + public void onFailure(InetAddressAndPort from, RequestFailure failure) { if (callback.invokeOnFailure()) callback.onFailure(from, failure); } diff --git a/test/unit/org/apache/cassandra/repair/FuzzTestBase.java b/test/unit/org/apache/cassandra/repair/FuzzTestBase.java index 9169fb4e88..593a494128 100644 --- a/test/unit/org/apache/cassandra/repair/FuzzTestBase.java +++ b/test/unit/org/apache/cassandra/repair/FuzzTestBase.java @@ -85,7 +85,7 @@ import org.apache.cassandra.dht.Murmur3Partitioner; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; import org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.gms.ApplicationState; import org.apache.cassandra.gms.EndpointState; import org.apache.cassandra.gms.HeartBeatState; @@ -816,7 +816,7 @@ public abstract class FuzzTestBase extends CQLTester.InMemory callback.onResponse(msg); } - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failureReason) { if (callback.invokeOnFailure()) callback.onFailure(from, failureReason); } @@ -950,7 +950,7 @@ public abstract class FuzzTestBase extends CQLTester.InMemory assert ctx == cb; try { - ctx.onFailure(to, RequestFailureReason.TIMEOUT); + ctx.onFailure(to, RequestFailure.TIMEOUT); } catch (Throwable t) { @@ -992,7 +992,7 @@ public abstract class FuzzTestBase extends CQLTester.InMemory } @Override - public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + public void onFailure(InetAddressAndPort from, RequestFailure failureReason) { promise.tryFailure(new MessagingService.FailureResponseException(from, failureReason)); } @@ -1177,7 +1177,7 @@ public abstract class FuzzTestBase extends CQLTester.InMemory try { if (msg.isFailureResponse()) - callback.onFailure(msg.from(), (RequestFailureReason) msg.payload); + callback.onFailure(msg.from(), (RequestFailure) msg.payload); else callback.onResponse(msg); } catch (Throwable t) diff --git a/test/unit/org/apache/cassandra/repair/messages/RepairMessageTest.java b/test/unit/org/apache/cassandra/repair/messages/RepairMessageTest.java index fb3ce470f5..7bbea89f12 100644 --- a/test/unit/org/apache/cassandra/repair/messages/RepairMessageTest.java +++ b/test/unit/org/apache/cassandra/repair/messages/RepairMessageTest.java @@ -23,7 +23,7 @@ import org.junit.Test; import org.apache.cassandra.concurrent.ScheduledExecutorPlus; import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.gms.IGossiper; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.metrics.RepairMetrics; @@ -87,7 +87,7 @@ public class RepairMessageTest public void noRetriesRequestFailed() { test(NO_RETRY_ATTEMPTS, ((ignore, callback) -> { - callback.onFailure(ADDRESS, RequestFailureReason.UNKNOWN); + callback.onFailure(ADDRESS, RequestFailure.UNKNOWN); assertNoRetries(); })); } @@ -105,7 +105,7 @@ public class RepairMessageTest public void retryWithTimeout() { test((maxAttempts, callback) -> { - callback.onFailure(ADDRESS, RequestFailureReason.TIMEOUT); + callback.onFailure(ADDRESS, RequestFailure.TIMEOUT); assertMetrics(maxAttempts, true, false); }); } @@ -114,7 +114,7 @@ public class RepairMessageTest public void retryWithFailure() { test((maxAttempts, callback) -> { - callback.onFailure(ADDRESS, RequestFailureReason.UNKNOWN); + callback.onFailure(ADDRESS, RequestFailure.UNKNOWN); assertMetrics(maxAttempts, false, true); }); } @@ -208,7 +208,7 @@ public class RepairMessageTest sendMessageWithRetries(ctx, backoff(maxAttempts), always(), PAYLOAD, VERB, ADDRESS, RepairMessage.NOOP_CALLBACK); for (int i = 0; i < maxAttempts; i++) - callback(messaging).onFailure(ADDRESS, RequestFailureReason.TIMEOUT); + callback(messaging).onFailure(ADDRESS, RequestFailure.TIMEOUT); fn.test(maxAttempts, callback(messaging)); Mockito.verifyNoInteractions(messaging); } diff --git a/test/unit/org/apache/cassandra/service/WriteResponseHandlerTest.java b/test/unit/org/apache/cassandra/service/WriteResponseHandlerTest.java index 03785f3c30..63562b35eb 100644 --- a/test/unit/org/apache/cassandra/service/WriteResponseHandlerTest.java +++ b/test/unit/org/apache/cassandra/service/WriteResponseHandlerTest.java @@ -35,8 +35,8 @@ import org.apache.cassandra.db.ConsistencyLevel; import org.apache.cassandra.db.Keyspace; import org.apache.cassandra.db.WriteType; import org.apache.cassandra.dht.Murmur3Partitioner; -import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.locator.EndpointsForToken; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.locator.NodeProximity; @@ -238,8 +238,8 @@ public class WriteResponseHandlerTest //Fail in remote DC - awr.onFailure(targets.get(3).endpoint(), RequestFailureReason.TIMEOUT); - awr.onFailure(targets.get(4).endpoint(), RequestFailureReason.TIMEOUT); + awr.onFailure(targets.get(3).endpoint(), RequestFailure.TIMEOUT); + awr.onFailure(targets.get(4).endpoint(), RequestFailure.TIMEOUT); awr.onResponse(createDummyMessage(5)); assertEquals(startingCountForWriteFailedIdealCL + 1, ks.metric.writeFailedIdealCL.getCount()); @@ -281,14 +281,14 @@ public class WriteResponseHandlerTest //Fail in local DC - awr.onFailure(targets.get(0).endpoint(), RequestFailureReason.TIMEOUT); - awr.onFailure(targets.get(1).endpoint(), RequestFailureReason.TIMEOUT); + awr.onFailure(targets.get(0).endpoint(), RequestFailure.TIMEOUT); + awr.onFailure(targets.get(1).endpoint(), RequestFailure.TIMEOUT); awr.onResponse(createDummyMessage(2)); //Fail in remote DC - awr.onFailure(targets.get(3).endpoint(), RequestFailureReason.TIMEOUT); - awr.onFailure(targets.get(4).endpoint(), RequestFailureReason.TIMEOUT); + awr.onFailure(targets.get(3).endpoint(), RequestFailure.TIMEOUT); + awr.onFailure(targets.get(4).endpoint(), RequestFailure.TIMEOUT); awr.onResponse(createDummyMessage(5)); assertEquals(startingCountForWriteFailedIdealCL, ks.metric.writeFailedIdealCL.getCount()); diff --git a/test/unit/org/apache/cassandra/service/accord/AccordJournalTest.java b/test/unit/org/apache/cassandra/service/accord/AccordJournalTest.java index 7059a12f2a..b24424b055 100644 --- a/test/unit/org/apache/cassandra/service/accord/AccordJournalTest.java +++ b/test/unit/org/apache/cassandra/service/accord/AccordJournalTest.java @@ -23,18 +23,21 @@ import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.List; +import org.junit.BeforeClass; import org.junit.Test; import accord.primitives.TxnId; import accord.utils.AccordGens; import accord.utils.Gen; import accord.utils.Gens; +import org.apache.cassandra.config.CassandraRelevantProperties; import org.apache.cassandra.io.util.DataInputBuffer; import org.apache.cassandra.io.util.DataOutputBuffer; import org.apache.cassandra.service.accord.AccordJournal.Key; import org.apache.cassandra.utils.AsymmetricOrdering; import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.FBUtilities.Order; +import org.apache.cassandra.utils.StorageCompatibilityMode; import org.checkerframework.checker.nullness.qual.Nullable; import static accord.utils.Property.qt; @@ -42,6 +45,12 @@ import static org.assertj.core.api.Assertions.assertThat; public class AccordJournalTest { + @BeforeClass + public static void setCompatibilityMode() + { + CassandraRelevantProperties.TEST_STORAGE_COMPATIBILITY_MODE.setEnum(StorageCompatibilityMode.NONE); + } + @Test public void keySerde() { diff --git a/test/unit/org/apache/cassandra/service/accord/AccordSyncPropagatorTest.java b/test/unit/org/apache/cassandra/service/accord/AccordSyncPropagatorTest.java index ab6e2790ce..5115dfdc47 100644 --- a/test/unit/org/apache/cassandra/service/accord/AccordSyncPropagatorTest.java +++ b/test/unit/org/apache/cassandra/service/accord/AccordSyncPropagatorTest.java @@ -59,7 +59,7 @@ import org.apache.cassandra.concurrent.AdaptingScheduledExecutorPlus; import org.apache.cassandra.concurrent.ScheduledExecutorPlus; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.dht.Murmur3Partitioner; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.gms.IFailureDetectionEventListener; import org.apache.cassandra.gms.IFailureDetector; import org.apache.cassandra.locator.InetAddressAndPort; @@ -289,10 +289,10 @@ public class AccordSyncPropagatorTest switch (action) { case ERROR: - cb.onFailure(to, RequestFailureReason.UNKNOWN); + cb.onFailure(to, RequestFailure.UNKNOWN); return; case TIMEOUT: - cb.onFailure(to, RequestFailureReason.TIMEOUT); + cb.onFailure(to, RequestFailure.TIMEOUT); return; case DELIVER: break; @@ -304,7 +304,7 @@ public class AccordSyncPropagatorTest scheduler.schedule(() -> { RequestCallback removed = callbacks.remove(message.id()); if (removed != null) - removed.onFailure(to, RequestFailureReason.TIMEOUT); + removed.onFailure(to, RequestFailure.TIMEOUT); }, 1, TimeUnit.MINUTES); } diff --git a/test/unit/org/apache/cassandra/service/paxos/PaxosVerbHandlerOutOfRangeTest.java b/test/unit/org/apache/cassandra/service/paxos/PaxosVerbHandlerOutOfRangeTest.java index 1ed2bef52e..a788c0ff9d 100644 --- a/test/unit/org/apache/cassandra/service/paxos/PaxosVerbHandlerOutOfRangeTest.java +++ b/test/unit/org/apache/cassandra/service/paxos/PaxosVerbHandlerOutOfRangeTest.java @@ -33,7 +33,7 @@ import org.apache.cassandra.ServerTestUtils; import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.db.Keyspace; import org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.metrics.StorageMetrics; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessagingService; @@ -48,9 +48,15 @@ import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.tcm.membership.NodeState; import org.apache.cassandra.utils.ByteBufferUtil; -import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.*; import static org.junit.Assert.assertEquals; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.MessageDelivery; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.broadcastAddress; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.bytesToken; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.node1; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.randomInt; +import static org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper.registerOutgoingMessageSink; + public class PaxosVerbHandlerOutOfRangeTest // PaxosV1 out of range tests - V2 implements OOTR checks at the protocol level { // For the purposes of this testing, the details of the Commit don't really matter @@ -175,7 +181,7 @@ public class PaxosVerbHandlerOutOfRangeTest // PaxosV1 out of range tests - V2 i MessageDelivery response = messageSink.get(100, TimeUnit.MILLISECONDS); assertEquals(verb, response.message.verb()); Assert.assertEquals(broadcastAddress, response.message.from()); - assertEquals(isOutOfRange, response.message.payload instanceof RequestFailureReason); + assertEquals(isOutOfRange, response.message.payload instanceof RequestFailure); assertEquals(messageId, response.message.id()); Assert.assertEquals(node1, response.to); assertEquals(startingTotalMetricCount + (isOutOfRange ? 1 : 0), StorageMetrics.totalOpsForInvalidToken.getCount()); diff --git a/test/unit/org/apache/cassandra/service/reads/ReadExecutorTest.java b/test/unit/org/apache/cassandra/service/reads/ReadExecutorTest.java index e23c7078b4..989130adb4 100644 --- a/test/unit/org/apache/cassandra/service/reads/ReadExecutorTest.java +++ b/test/unit/org/apache/cassandra/service/reads/ReadExecutorTest.java @@ -22,32 +22,33 @@ import java.util.concurrent.TimeUnit; import org.apache.commons.lang3.exception.ExceptionUtils; -import org.apache.cassandra.ServerTestUtils; -import org.apache.cassandra.dht.Murmur3Partitioner; -import org.apache.cassandra.dht.Token; -import org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper; -import org.apache.cassandra.locator.ReplicaPlan; import org.junit.Before; import org.junit.BeforeClass; import org.junit.Test; import org.apache.cassandra.SchemaLoader; +import org.apache.cassandra.ServerTestUtils; import org.apache.cassandra.Util; import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.db.ConsistencyLevel; import org.apache.cassandra.db.Keyspace; import org.apache.cassandra.db.SinglePartitionReadCommand; +import org.apache.cassandra.dht.Murmur3Partitioner; +import org.apache.cassandra.dht.Token; +import org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper; import org.apache.cassandra.exceptions.ReadFailureException; import org.apache.cassandra.exceptions.ReadTimeoutException; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.locator.EndpointsForToken; import org.apache.cassandra.locator.InetAddressAndPort; -import org.apache.cassandra.tcm.Epoch; +import org.apache.cassandra.locator.ReplicaPlan; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.NoPayload; import org.apache.cassandra.net.Verb; import org.apache.cassandra.schema.KeyspaceParams; import org.apache.cassandra.transport.Dispatcher; +import org.apache.cassandra.tcm.Epoch; import static java.util.concurrent.TimeUnit.DAYS; import static java.util.concurrent.TimeUnit.MILLISECONDS; @@ -149,8 +150,8 @@ public class ReadExecutorTest public void run() { //Failures end the read promptly but don't require mock data to be suppleid - executor.handler.onFailure(targets.get(0).endpoint(), RequestFailureReason.READ_TOO_MANY_TOMBSTONES); - executor.handler.onFailure(targets.get(1).endpoint(), RequestFailureReason.READ_TOO_MANY_TOMBSTONES); + executor.handler.onFailure(targets.get(0).endpoint(), RequestFailure.READ_TOO_MANY_TOMBSTONES); + executor.handler.onFailure(targets.get(1).endpoint(), RequestFailure.READ_TOO_MANY_TOMBSTONES); executor.handler.condition.signalAll(); } }.start(); @@ -221,7 +222,7 @@ public class ReadExecutorTest { // Fail the first request. When this fails the number of contacts has already been increased // to 2, so the failure won't actally signal. However... - executor.handler.onFailure(targets.get(0).endpoint(), RequestFailureReason.READ_TOO_MANY_TOMBSTONES); + executor.handler.onFailure(targets.get(0).endpoint(), RequestFailure.READ_TOO_MANY_TOMBSTONES); // ...speculative retries are fired after a short wait, and it is possible for the failure to // reach the handler just before one is fired and the number of contacts incremented... diff --git a/test/unit/org/apache/cassandra/tcm/DiscoverySimulationTest.java b/test/unit/org/apache/cassandra/tcm/DiscoverySimulationTest.java index 564eb98f84..6021355442 100644 --- a/test/unit/org/apache/cassandra/tcm/DiscoverySimulationTest.java +++ b/test/unit/org/apache/cassandra/tcm/DiscoverySimulationTest.java @@ -36,7 +36,7 @@ import org.junit.Test; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.dht.Murmur3Partitioner; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.net.ConnectionType; import org.apache.cassandra.net.IVerbHandler; @@ -194,7 +194,7 @@ public class DiscoverySimulationTest else { logger.info("{} simulating failure sending request to {}", addr, to); - cb.onFailure(to, RequestFailureReason.TIMEOUT); + cb.onFailure(to, RequestFailure.TIMEOUT); } } catch (IOException e) diff --git a/test/unit/org/apache/cassandra/tcm/sequences/ProgressBarrierTest.java b/test/unit/org/apache/cassandra/tcm/sequences/ProgressBarrierTest.java index ace1da7f1a..1006bc941f 100644 --- a/test/unit/org/apache/cassandra/tcm/sequences/ProgressBarrierTest.java +++ b/test/unit/org/apache/cassandra/tcm/sequences/ProgressBarrierTest.java @@ -37,7 +37,7 @@ import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.ConsistencyLevel; import org.apache.cassandra.distributed.api.IIsolatedExecutor; import org.apache.cassandra.distributed.test.log.CMSTestBase; -import org.apache.cassandra.exceptions.RequestFailureReason; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.harry.gen.EntropySource; import org.apache.cassandra.harry.gen.Surjections; import org.apache.cassandra.harry.gen.rng.PCGFastPure; @@ -147,7 +147,7 @@ public class ProgressBarrierTest extends CMSTestBase } else { - cb.onFailure(message.from(), RequestFailureReason.TIMEOUT); + cb.onFailure(message.from(), RequestFailure.TIMEOUT); } }