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
This commit is contained in:
Ariel Weisberg 2023-08-18 16:48:39 -04:00 committed by David Capwell
parent bfa0e59f7f
commit 82acd3e950
64 changed files with 960 additions and 278 deletions

@ -1 +1 @@
Subproject commit 2ad55e03c43ce074cdf5e36cfa14cb4278c2dc0f
Subproject commit 91336705bde8332954e849219d73205d68fa168a

View File

@ -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<StackTraceElement> stackTraceElementSerializer = new IVersionedSerializer<StackTraceElement>()
{
@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<Throwable> remoteExceptionSerializer = new IVersionedSerializer<Throwable>()
{
@Override
public void serialize(Throwable t, DataOutputPlus out, int version) throws IOException
{
Map<Throwable, Integer> alreadySerialized = new IdentityHashMap<>();
serializeNextException(t, out, true, version, 0, alreadySerialized);
}
private int serializeNextException(Throwable t, DataOutputPlus out, boolean isFirstException, int version, int nextExceptionId, Map<Throwable, Integer> 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<Integer, Throwable> alreadyDeserialized = new HashMap<>();
return deserializeNextException(in, version, alreadyDeserialized);
}
private Throwable deserializeNextException(DataInputPlus in, int version, Map<Integer, Throwable> 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<Throwable> alreadySeen = newSetFromMap(new IdentityHashMap<>());
return nextExceptionSerializedSize(t, version, true, alreadySeen);
}
private long nextExceptionSerializedSize(Throwable t, int version, boolean isFirstException, Set<Throwable> 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<Throwable> nullableRemoteExceptionSerializer = NullableSerializer.wrap(remoteExceptionSerializer);
}

View File

@ -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<RequestFailure> serializer = new IVersionedSerializer<RequestFailure>()
{
@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 + '\'' +
'}';
}
}

View File

@ -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;

View File

@ -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();

View File

@ -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<E extends Endpoints<E>, P extends ReplicaPlan<E, P>>
{
Epoch epoch();
@ -49,7 +50,7 @@ public interface ReplicaPlan<E extends Endpoints<E>, P extends ReplicaPlan<E, P>
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<E extends Endpoints<E>, P extends ReplicaPlan.ForRead<E, P>> extends ReplicaPlan<E, P>
@ -115,7 +116,7 @@ public interface ReplicaPlan<E extends Endpoints<E>, P extends ReplicaPlan<E, P>
contacted.add(addr);
}
public void collectFailure(InetAddressAndPort inetAddressAndPort, RequestFailureReason t) {}
public void collectFailure(InetAddressAndPort inetAddressAndPort, RequestFailure t) {}
}

View File

@ -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<RequestFailureReason> response = Message.failureResponse(header.id,
header.expiresAtNanos,
RequestFailureReason.forException(failure));
Message<RequestFailure> response = Message.failureResponse(header.id,
header.expiresAtNanos,
RequestFailure.forException(failure));
messaging.send(response, to);
}
}

View File

@ -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<T> implements ReplyContext
}
/** Builds a failure response Message with an explicit reason, and fields inferred from request Message */
public Message<RequestFailureReason> failureResponse(RequestFailureReason reason)
public Message<RequestFailure> failureResponse(RequestFailureReason reason)
{
return failureResponse(id(), expiresAtNanos(), reason);
return failureResponse(reason, null);
}
static Message<RequestFailureReason> failureResponse(long id, long expiresAtNanos, RequestFailureReason reason)
public Message<RequestFailure> failureResponse(RequestFailureReason reason, @Nullable Throwable failure)
{
return failureResponse(id(), expiresAtNanos(), new RequestFailure(reason, failure));
}
static Message<RequestFailure> failureResponse(long id, long expiresAtNanos, RequestFailure reason)
{
return outWithParam(id, Verb.FAILURE_RSP, expiresAtNanos, reason, null, null);
}

View File

@ -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 <V> 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 <REQ, RSP> 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;
}

View File

@ -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());

View File

@ -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<T>
/**
* 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)
{
}

View File

@ -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<T> extends RequestCallback<T>
@ -26,7 +26,7 @@ public interface RequestCallbackWithFailure<T> extends RequestCallback<T>
/**
* 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

View File

@ -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)

View File

@ -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
{

View File

@ -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

View File

@ -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<InetAddressAndPort> 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));
}
}
}

View File

@ -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

View File

@ -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<T> 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<T> implements RequestCallback
if (failureReasonByEndpoint == null)
failureReasonByEndpoint = new ConcurrentHashMap<>();
}
failureReasonByEndpoint.put(from, failureReason);
failureReasonByEndpoint.put(from, failure.reason);
logFailureOrTimeoutToIdealCLDelegate();

View File

@ -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 " +

View File

@ -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<T> extends AbstractWriteResponseHandler<T>
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()

View File

@ -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<T> implements RequestCallbackWith
private static final AtomicReferenceFieldUpdater<FailureRecordingCallback, FailureResponses> 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)

View File

@ -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;
}
}

View File

@ -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<TruncateResponse
}
@Override
public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason)
public void onFailure(InetAddressAndPort from, RequestFailure failure)
{
// If the truncation hasn't succeeded on some replica, abort and indicate this back to the client.
failureReasonByEndpoint.put(from, failureReason);
failureReasonByEndpoint.put(from, failure.reason);
condition.signalAll();
}

View File

@ -24,8 +24,9 @@ import org.slf4j.LoggerFactory;
import accord.coordinate.Timeout;
import accord.local.AgentExecutor;
import accord.messages.Callback;
import accord.messages.SafeCallback;
import accord.messages.Reply;
import accord.messages.SafeCallback;
import org.apache.cassandra.exceptions.RequestFailure;
import org.apache.cassandra.exceptions.RequestFailureReason;
import org.apache.cassandra.locator.InetAddressAndPort;
import org.apache.cassandra.net.Message;
@ -49,19 +50,19 @@ class AccordCallback<T extends Reply> extends SafeCallback<T> 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

View File

@ -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());

View File

@ -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);
}

View File

@ -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)
{
}

View File

@ -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<OnDone extends Consumer<? super PaxosCommit.Status>> 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<OnDone extends Consumer<? super PaxosCommit.Status>> 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);

View File

@ -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<PaxosPrepare.Response> 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<PaxosPrepare.Response> 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<PaxosPrepare.Response> im
}
@Override
public void onRefreshFailure(InetAddressAndPort from, RequestFailureReason reason)
public void onRefreshFailure(InetAddressAndPort from, RequestFailure reason)
{
onFailure(from, reason);
}

View File

@ -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<PaxosPrep
interface Callbacks
{
void onRefreshFailure(InetAddressAndPort from, RequestFailureReason reason);
void onRefreshFailure(InetAddressAndPort from, RequestFailure reason);
void onRefreshSuccess(Ballot isSupersededBy, InetAddressAndPort from);
}
@ -102,7 +101,7 @@ public class PaxosPrepareRefresh implements RequestCallbackWithFailure<PaxosPrep
}
@Override
public void onFailure(InetAddressAndPort from, RequestFailureReason reason)
public void onFailure(InetAddressAndPort from, RequestFailure reason)
{
callbacks.onRefreshFailure(from, reason);
}
@ -124,8 +123,8 @@ public class PaxosPrepareRefresh implements RequestCallbackWithFailure<PaxosPrep
}
catch (Exception ex)
{
RequestFailureReason reason = UNKNOWN;
if (ex instanceof WriteTimeoutException) reason = TIMEOUT;
RequestFailure reason = RequestFailure.UNKNOWN;
if (ex instanceof WriteTimeoutException) reason = RequestFailure.TIMEOUT;
else logger.error("Failed to apply paxos refresh-prepare locally", ex);
onFailure(getBroadcastAddressAndPort(), reason);
@ -167,7 +166,7 @@ public class PaxosPrepareRefresh implements RequestCallbackWithFailure<PaxosPrep
{
Response response = execute(message.payload, message.from());
if (response == null)
MessagingService.instance().respondWithFailure(UNKNOWN, message);
MessagingService.instance().respondWithFailure(RequestFailureReason.UNKNOWN, message);
else
MessagingService.instance().respond(response, message);
}

View File

@ -29,7 +29,7 @@ import org.slf4j.LoggerFactory;
import org.apache.cassandra.db.Keyspace;
import org.apache.cassandra.db.TypeSizes;
import org.apache.cassandra.exceptions.RequestFailureReason;
import org.apache.cassandra.exceptions.RequestFailure;
import org.apache.cassandra.io.IVersionedSerializer;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus;
@ -43,8 +43,8 @@ import org.apache.cassandra.utils.concurrent.ConditionAsConsumer;
import static java.util.Collections.emptyMap;
import static org.apache.cassandra.exceptions.RequestFailureReason.UNKNOWN;
import static org.apache.cassandra.net.Verb.PAXOS2_PROPOSE_REQ;
import static org.apache.cassandra.service.paxos.PaxosPropose.Superseded.SideEffects.NO;
import static org.apache.cassandra.service.paxos.PaxosPropose.Superseded.SideEffects.MAYBE;
import static org.apache.cassandra.service.paxos.PaxosPropose.Superseded.SideEffects.NO;
import static org.apache.cassandra.utils.Clock.Global.nanoTime;
import static org.apache.cassandra.utils.concurrent.ConditionAsConsumer.newConditionAsConsumer;
@ -263,7 +263,7 @@ public class PaxosPropose<OnDone extends Consumer<? super PaxosPropose.Status>>
}
@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);

View File

@ -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());
}

View File

@ -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<T> extends FailureRecordingCallback<T>
@ -58,7 +58,7 @@ public abstract class PaxosRequestCallback<T> extends FailureRecordingCallback<T
}
catch (Exception ex)
{
RequestFailureReason reason = UNKNOWN;
RequestFailure reason = UNKNOWN;
if (ex instanceof WriteTimeoutException) reason = TIMEOUT;
else logger.error("Failed to apply {} locally", parameter, ex);

View File

@ -19,7 +19,12 @@
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.Collections;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.TypeSizes;
@ -27,7 +32,7 @@ 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.io.IVersionedSerializer;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus;
@ -78,7 +83,7 @@ public class PaxosCleanupComplete extends AsyncFuture<Void> 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));
}

View File

@ -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<Void> 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);
}

View File

@ -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<Void> 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));
}

View File

@ -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<PaxosCleanupHistory> 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));
}

View File

@ -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<E extends Endpoints<E>, 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<E extends Endpoints<E>, 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();

View File

@ -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<ParamType, Object> params, InetAddressAndPort from)
public RequestFailure updateCounters(Map<ParamType, Object> params, InetAddressAndPort from)
{
for (Map.Entry<ParamType, Object> 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;

View File

@ -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()

View File

@ -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));

View File

@ -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;

View File

@ -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()

View File

@ -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<SystemInfo> systemInfoSupplier = Suppliers.memoize(SystemInfo::new);
public static void setAvailableProcessors(int value)

View File

@ -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"

View File

@ -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;

View File

@ -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();
/**

View File

@ -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);
}

View File

@ -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());

View File

@ -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());

View File

@ -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<Throwable, Throwable> 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;
}
}

View File

@ -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();
}

View File

@ -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();

View File

@ -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<NoPayload> 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<RequestFailureReason> 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<RequestFailure> 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<NoPayload> 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);
}

View File

@ -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);
}

View File

@ -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)

View File

@ -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);
}

View File

@ -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());

View File

@ -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()
{

View File

@ -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);
}

View File

@ -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());

View File

@ -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...

View File

@ -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)

View File

@ -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);
}
}