Repository: cassandra Updated Branches: refs/heads/trunk 7b61b0be8 -> a09dc3a53
Clean up Message.Request implementations patch by Aleksey Yeschenko; reviewed by Chris Lohfink for CASSANDRA-14677 Project: http://git-wip-us.apache.org/repos/asf/cassandra/repo Commit: http://git-wip-us.apache.org/repos/asf/cassandra/commit/a09dc3a5 Tree: http://git-wip-us.apache.org/repos/asf/cassandra/tree/a09dc3a5 Diff: http://git-wip-us.apache.org/repos/asf/cassandra/diff/a09dc3a5 Branch: refs/heads/trunk Commit: a09dc3a530921c5c4a8a9cdf3df201aef2c11742 Parents: 7b61b0b Author: Aleksey Yeshchenko <[email protected]> Authored: Wed Aug 29 18:42:39 2018 +0100 Committer: Aleksey Yeshchenko <[email protected]> Committed: Fri Aug 31 10:36:46 2018 +0100 ---------------------------------------------------------------------- CHANGES.txt | 1 + .../apache/cassandra/audit/AuditLogEntry.java | 25 ++- .../apache/cassandra/audit/AuditLogManager.java | 4 +- .../apache/cassandra/cql3/QueryProcessor.java | 6 +- .../apache/cassandra/service/ClientState.java | 6 + .../apache/cassandra/service/QueryState.java | 57 ++---- .../cassandra/service/StorageService.java | 5 + .../org/apache/cassandra/transport/Message.java | 68 +++++-- .../cassandra/transport/ServerConnection.java | 32 +--- .../transport/messages/AuthResponse.java | 45 +++-- .../transport/messages/BatchMessage.java | 91 +++++----- .../transport/messages/ExecuteMessage.java | 177 +++++++++---------- .../transport/messages/OptionsMessage.java | 3 +- .../transport/messages/PrepareMessage.java | 87 +++++---- .../transport/messages/QueryMessage.java | 114 ++++++------ .../transport/messages/RegisterMessage.java | 3 +- .../transport/messages/StartupMessage.java | 3 +- 17 files changed, 378 insertions(+), 349 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/cassandra/blob/a09dc3a5/CHANGES.txt ---------------------------------------------------------------------- diff --git a/CHANGES.txt b/CHANGES.txt index 249e034..6489038 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 4.0 + * Clean up Message.Request implementations (CASSANDRA-14677) * Disable old native protocol versions on demand (CASANDRA-14659) * Allow specifying now-in-seconds in native protocol (CASSANDRA-14664) * Improve BTree build performance by avoiding data copy (CASSANDRA-9989) http://git-wip-us.apache.org/repos/asf/cassandra/blob/a09dc3a5/src/java/org/apache/cassandra/audit/AuditLogEntry.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/audit/AuditLogEntry.java b/src/java/org/apache/cassandra/audit/AuditLogEntry.java index d53fc6e..0b891d4 100644 --- a/src/java/org/apache/cassandra/audit/AuditLogEntry.java +++ b/src/java/org/apache/cassandra/audit/AuditLogEntry.java @@ -44,8 +44,18 @@ public class AuditLogEntry private final String scope; private final String operation; private final QueryOptions options; - - private AuditLogEntry(AuditLogEntryType type, InetAddressAndPort source, String user, long timestamp, UUID batch, String keyspace, String scope, String operation, QueryOptions options) + private final QueryState state; + + private AuditLogEntry(AuditLogEntryType type, + InetAddressAndPort source, + String user, + long timestamp, + UUID batch, + String keyspace, + String scope, + String operation, + QueryOptions options, + QueryState state) { this.type = type; this.source = source; @@ -56,6 +66,7 @@ public class AuditLogEntry this.scope = scope; this.operation = operation; this.options = options; + this.state = state; } String getLogString() @@ -170,9 +181,14 @@ public class AuditLogEntry private String scope; private String operation; private QueryOptions options; + private QueryState state; - public Builder(ClientState clientState) + public Builder(QueryState queryState) { + state = queryState; + + ClientState clientState = queryState.getClientState(); + if (clientState != null) { if (clientState.getRemoteAddress() != null) @@ -207,6 +223,7 @@ public class AuditLogEntry scope = entry.scope; operation = entry.operation; options = entry.options; + state = entry.state; } public Builder setType(AuditLogEntryType type) @@ -291,7 +308,7 @@ public class AuditLogEntry public AuditLogEntry build() { timestamp = timestamp > 0 ? timestamp : System.currentTimeMillis(); - return new AuditLogEntry(type, source, user, timestamp, batch, keyspace, scope, operation, options); + return new AuditLogEntry(type, source, user, timestamp, batch, keyspace, scope, operation, options, state); } } } http://git-wip-us.apache.org/repos/asf/cassandra/blob/a09dc3a5/src/java/org/apache/cassandra/audit/AuditLogManager.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/audit/AuditLogManager.java b/src/java/org/apache/cassandra/audit/AuditLogManager.java index 9e6a0a1..ab9c2e9 100644 --- a/src/java/org/apache/cassandra/audit/AuditLogManager.java +++ b/src/java/org/apache/cassandra/audit/AuditLogManager.java @@ -214,7 +214,7 @@ public class AuditLogManager List<AuditLogEntry> auditLogEntries = new ArrayList<>(queryOrIdList.size() + 1); UUID batchId = UUID.randomUUID(); String queryString = String.format("BatchId:[%s] - BATCH of [%d] statements", batchId, queryOrIdList.size()); - AuditLogEntry entry = new AuditLogEntry.Builder(state.getClientState()) + AuditLogEntry entry = new AuditLogEntry.Builder(state) .setOperation(queryString) .setOptions(options) .setTimestamp(queryStartTimeMillis) @@ -226,7 +226,7 @@ public class AuditLogManager for (int i = 0; i < queryOrIdList.size(); i++) { CQLStatement statement = prepared.get(i).statement; - entry = new AuditLogEntry.Builder(state.getClientState()) + entry = new AuditLogEntry.Builder(state) .setType(statement.getAuditLogContext().auditLogEntryType) .setOperation(prepared.get(i).rawCQLStatement) .setTimestamp(queryStartTimeMillis) http://git-wip-us.apache.org/repos/asf/cassandra/blob/a09dc3a5/src/java/org/apache/cassandra/cql3/QueryProcessor.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/cql3/QueryProcessor.java b/src/java/org/apache/cassandra/cql3/QueryProcessor.java index 77b4cdc..79e19c1 100644 --- a/src/java/org/apache/cassandra/cql3/QueryProcessor.java +++ b/src/java/org/apache/cassandra/cql3/QueryProcessor.java @@ -122,11 +122,11 @@ public class QueryProcessor implements QueryHandler { INSTANCE; - private final QueryState queryState; + private final ClientState clientState; InternalStateInstance() { - queryState = new QueryState(ClientState.forInternalCalls(SchemaConstants.SYSTEM_KEYSPACE_NAME)); + clientState = ClientState.forInternalCalls(SchemaConstants.SYSTEM_KEYSPACE_NAME); } } @@ -164,7 +164,7 @@ public class QueryProcessor implements QueryHandler private static QueryState internalQueryState() { - return InternalStateInstance.INSTANCE.queryState; + return new QueryState(InternalStateInstance.INSTANCE.clientState); } private QueryProcessor() http://git-wip-us.apache.org/repos/asf/cassandra/blob/a09dc3a5/src/java/org/apache/cassandra/service/ClientState.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/service/ClientState.java b/src/java/org/apache/cassandra/service/ClientState.java index 688df91..cb06161 100644 --- a/src/java/org/apache/cassandra/service/ClientState.java +++ b/src/java/org/apache/cassandra/service/ClientState.java @@ -17,6 +17,7 @@ */ package org.apache.cassandra.service; +import java.net.InetAddress; import java.net.InetSocketAddress; import java.net.SocketAddress; import java.util.Arrays; @@ -291,6 +292,11 @@ public class ClientState return remoteAddress; } + InetAddress getClientAddress() + { + return isInternal ? null : remoteAddress.getAddress(); + } + public String getRawKeyspace() { return keyspace; http://git-wip-us.apache.org/repos/asf/cassandra/blob/a09dc3a5/src/java/org/apache/cassandra/service/QueryState.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/service/QueryState.java b/src/java/org/apache/cassandra/service/QueryState.java index f0ae3b2..b266fb8 100644 --- a/src/java/org/apache/cassandra/service/QueryState.java +++ b/src/java/org/apache/cassandra/service/QueryState.java @@ -18,12 +18,8 @@ package org.apache.cassandra.service; import java.net.InetAddress; -import java.nio.ByteBuffer; -import java.util.Map; -import java.util.UUID; -import java.util.concurrent.ThreadLocalRandom; -import org.apache.cassandra.tracing.Tracing; +import org.apache.cassandra.utils.FBUtilities; /** * Represents the state related to a given query. @@ -31,7 +27,9 @@ import org.apache.cassandra.tracing.Tracing; public class QueryState { private final ClientState clientState; - private volatile UUID preparedTracingSession; + + private long timestamp = Long.MIN_VALUE; + private int nowInSeconds = Integer.MIN_VALUE; public QueryState(ClientState clientState) { @@ -51,49 +49,22 @@ public class QueryState return clientState; } - /** - * This clock guarantees that updates for the same QueryState will be ordered - * in the sequence seen, even if multiple updates happen in the same millisecond. - */ - public long getTimestamp() - { - return clientState.getTimestamp(); - } - - public boolean traceNextQuery() - { - if (preparedTracingSession != null) - { - return true; - } - - double traceProbability = StorageService.instance.getTraceProbability(); - return traceProbability != 0 && ThreadLocalRandom.current().nextDouble() < traceProbability; - } - - public void prepareTracingSession(UUID sessionId) + public InetAddress getClientAddress() { - this.preparedTracingSession = sessionId; + return clientState.getClientAddress(); } - public void createTracingSession(Map<String,ByteBuffer> customPayload) + public long getTimestamp() { - UUID session = this.preparedTracingSession; - if (session == null) - { - Tracing.instance.newSession(customPayload); - } - else - { - Tracing.instance.newSession(session, customPayload); - this.preparedTracingSession = null; - } + if (timestamp == Long.MIN_VALUE) + timestamp = clientState.getTimestamp(); + return timestamp; } - public InetAddress getClientAddress() + public int getNowInSeconds() { - return clientState.isInternal - ? null - : clientState.getRemoteAddress().getAddress(); + if (nowInSeconds == Integer.MIN_VALUE) + nowInSeconds = FBUtilities.nowInSeconds(); + return nowInSeconds; } } http://git-wip-us.apache.org/repos/asf/cassandra/blob/a09dc3a5/src/java/org/apache/cassandra/service/StorageService.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index a3c61a3..09bed8d 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -5352,6 +5352,11 @@ public class StorageService extends NotificationBroadcasterSupport implements IE return traceProbability; } + public boolean shouldTraceProbablistically() + { + return traceProbability != 0 && ThreadLocalRandom.current().nextDouble() < traceProbability; + } + public void disableAutoCompaction(String ks, String... tables) throws IOException { for (ColumnFamilyStore cfs : getValidColumnFamilies(true, true, ks, tables)) http://git-wip-us.apache.org/repos/asf/cassandra/blob/a09dc3a5/src/java/org/apache/cassandra/transport/Message.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/transport/Message.java b/src/java/org/apache/cassandra/transport/Message.java index 413c94a..255af0e 100644 --- a/src/java/org/apache/cassandra/transport/Message.java +++ b/src/java/org/apache/cassandra/transport/Message.java @@ -42,11 +42,13 @@ import com.google.common.collect.ImmutableSet; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.apache.cassandra.audit.AuditLogManager; import org.apache.cassandra.service.ClientWarn; +import org.apache.cassandra.service.StorageService; +import org.apache.cassandra.tracing.Tracing; import org.apache.cassandra.transport.messages.*; import org.apache.cassandra.service.QueryState; import org.apache.cassandra.utils.JVMStabilityInspector; +import org.apache.cassandra.utils.UUIDGen; /** * A message from the CQL binary protocol. @@ -201,11 +203,7 @@ public abstract class Message public static abstract class Request extends Message { - protected boolean tracingRequested; - - protected final AuditLogManager auditLogManager = AuditLogManager.getInstance(); - protected boolean auditLogEnabled = auditLogManager.isAuditingEnabled(); - protected boolean isLoggingEnabled = auditLogManager.isLoggingEnabled(); + private boolean tracingRequested; protected Request(Type type) { @@ -215,14 +213,56 @@ public abstract class Message throw new IllegalArgumentException(); } - public abstract Response execute(QueryState queryState, long queryStartNanoTime); + protected boolean isTraceable() + { + return false; + } + + protected abstract Response execute(QueryState queryState, long queryStartNanoTime, boolean traceRequest); + + final Response execute(QueryState queryState, long queryStartNanoTime) + { + boolean shouldTrace = false; + UUID tracingSessionId = null; + + if (isTraceable()) + { + if (isTracingRequested()) + { + shouldTrace = true; + tracingSessionId = UUIDGen.getTimeUUID(); + Tracing.instance.newSession(tracingSessionId, getCustomPayload()); + } + else if (StorageService.instance.shouldTraceProbablistically()) + { + shouldTrace = true; + Tracing.instance.newSession(getCustomPayload()); + } + } + + Response response; + try + { + response = execute(queryState, queryStartNanoTime, shouldTrace); + } + finally + { + if (shouldTrace) + Tracing.instance.stopSession(); + } + + if (isTraceable() && isTracingRequested()) + response.setTracingId(tracingSessionId); + + return response; + } - public void setTracingRequested() + void setTracingRequested() { - this.tracingRequested = true; + tracingRequested = true; } - public boolean isTracingRequested() + boolean isTracingRequested() { return tracingRequested; } @@ -241,18 +281,18 @@ public abstract class Message throw new IllegalArgumentException(); } - public Message setTracingId(UUID tracingId) + Message setTracingId(UUID tracingId) { this.tracingId = tracingId; return this; } - public UUID getTracingId() + UUID getTracingId() { return tracingId; } - public Message setWarnings(List<String> warnings) + Message setWarnings(List<String> warnings) { this.warnings = warnings; return this; @@ -565,7 +605,7 @@ public abstract class Message if (connection.getVersion().isGreaterOrEqualTo(ProtocolVersion.V4)) ClientWarn.instance.captureWarnings(); - QueryState qstate = connection.validateNewMessage(request.type, connection.getVersion(), request.getStreamId()); + QueryState qstate = connection.validateNewMessage(request.type, connection.getVersion()); logger.trace("Received: {}, v={}", request, connection.getVersion()); connection.requests.inc(); http://git-wip-us.apache.org/repos/asf/cassandra/blob/a09dc3a5/src/java/org/apache/cassandra/transport/ServerConnection.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/transport/ServerConnection.java b/src/java/org/apache/cassandra/transport/ServerConnection.java index 00e334c..de8a02a 100644 --- a/src/java/org/apache/cassandra/transport/ServerConnection.java +++ b/src/java/org/apache/cassandra/transport/ServerConnection.java @@ -17,9 +17,6 @@ */ package org.apache.cassandra.transport; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; - import javax.net.ssl.SSLPeerUnverifiedException; import javax.security.cert.X509Certificate; @@ -36,33 +33,19 @@ import org.apache.cassandra.service.QueryState; public class ServerConnection extends Connection { - private static Logger logger = LoggerFactory.getLogger(ServerConnection.class); + private static final Logger logger = LoggerFactory.getLogger(ServerConnection.class); private volatile IAuthenticator.SaslNegotiator saslNegotiator; private final ClientState clientState; private volatile ConnectionStage stage; public final Counter requests = new Counter(); - private final ConcurrentMap<Integer, QueryState> queryStates = new ConcurrentHashMap<>(); - - public ServerConnection(Channel channel, ProtocolVersion version, Connection.Tracker tracker) + ServerConnection(Channel channel, ProtocolVersion version, Connection.Tracker tracker) { super(channel, version, tracker); - this.clientState = ClientState.forExternalCalls(channel.remoteAddress()); - this.stage = ConnectionStage.ESTABLISHED; - } - private QueryState getQueryState(int streamId) - { - QueryState qState = queryStates.get(streamId); - if (qState == null) - { - // In theory we shouldn't get any race here, but it never hurts to be careful - QueryState newState = new QueryState(clientState); - if ((qState = queryStates.putIfAbsent(streamId, newState)) == null) - qState = newState; - } - return qState; + clientState = ClientState.forExternalCalls(channel.remoteAddress()); + stage = ConnectionStage.ESTABLISHED; } public ClientState getClientState() @@ -75,7 +58,7 @@ public class ServerConnection extends Connection return stage; } - public QueryState validateNewMessage(Message.Type type, ProtocolVersion version, int streamId) + QueryState validateNewMessage(Message.Type type, ProtocolVersion version) { switch (stage) { @@ -95,10 +78,11 @@ public class ServerConnection extends Connection default: throw new AssertionError(); } - return getQueryState(streamId); + + return new QueryState(clientState); } - public void applyStateTransition(Message.Type requestType, Message.Type responseType) + void applyStateTransition(Message.Type requestType, Message.Type responseType) { switch (stage) { http://git-wip-us.apache.org/repos/asf/cassandra/blob/a09dc3a5/src/java/org/apache/cassandra/transport/messages/AuthResponse.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/transport/messages/AuthResponse.java b/src/java/org/apache/cassandra/transport/messages/AuthResponse.java index 2a20898..d101681 100644 --- a/src/java/org/apache/cassandra/transport/messages/AuthResponse.java +++ b/src/java/org/apache/cassandra/transport/messages/AuthResponse.java @@ -22,6 +22,7 @@ import java.nio.ByteBuffer; import io.netty.buffer.ByteBuf; import org.apache.cassandra.audit.AuditLogEntry; import org.apache.cassandra.audit.AuditLogEntryType; +import org.apache.cassandra.audit.AuditLogManager; import org.apache.cassandra.auth.AuthenticatedUser; import org.apache.cassandra.auth.IAuthenticator; import org.apache.cassandra.exceptions.AuthenticationException; @@ -70,8 +71,10 @@ public class AuthResponse extends Message.Request } @Override - public Response execute(QueryState queryState, long queryStartNanoTime) + protected Response execute(QueryState queryState, long queryStartNanoTime, boolean traceRequest) { + AuditLogManager auditLogManager = AuditLogManager.getInstance(); + try { IAuthenticator.SaslNegotiator negotiator = ((ServerConnection) connection).getSaslNegotiator(queryState); @@ -81,14 +84,8 @@ public class AuthResponse extends Message.Request AuthenticatedUser user = negotiator.getAuthenticatedUser(); queryState.getClientState().login(user); ClientMetrics.instance.markAuthSuccess(); - if (auditLogEnabled) - { - AuditLogEntry auditEntry = new AuditLogEntry.Builder(queryState.getClientState()) - .setOperation("LOGIN SUCCESSFUL") - .setType(AuditLogEntryType.LOGIN_SUCCESS) - .build(); - auditLogManager.log(auditEntry); - } + if (auditLogManager.isAuditingEnabled()) + logSuccess(queryState); // authentication is complete, send a ready message to the client return new AuthSuccess(challenge); } @@ -100,15 +97,29 @@ public class AuthResponse extends Message.Request catch (AuthenticationException e) { ClientMetrics.instance.markAuthFailure(); - if (auditLogEnabled) - { - AuditLogEntry auditEntry = new AuditLogEntry.Builder(queryState.getClientState()) - .setOperation("LOGIN FAILURE") - .setType(AuditLogEntryType.LOGIN_ERROR) - .build(); - auditLogManager.log(auditEntry, e); - } + if (auditLogManager.isAuditingEnabled()) + logException(queryState, e); return ErrorMessage.fromException(e); } } + + private void logSuccess(QueryState state) + { + AuditLogEntry entry = + new AuditLogEntry.Builder(state) + .setOperation("LOGIN SUCCESSFUL") + .setType(AuditLogEntryType.LOGIN_SUCCESS) + .build(); + AuditLogManager.getInstance().log(entry); + } + + private void logException(QueryState state, AuthenticationException e) + { + AuditLogEntry entry = + new AuditLogEntry.Builder(state) + .setOperation("LOGIN FAILURE") + .setType(AuditLogEntryType.LOGIN_ERROR) + .build(); + AuditLogManager.getInstance().log(entry, e); + } } http://git-wip-us.apache.org/repos/asf/cassandra/blob/a09dc3a5/src/java/org/apache/cassandra/transport/messages/BatchMessage.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/transport/messages/BatchMessage.java b/src/java/org/apache/cassandra/transport/messages/BatchMessage.java index 29f92f7..b379918 100644 --- a/src/java/org/apache/cassandra/transport/messages/BatchMessage.java +++ b/src/java/org/apache/cassandra/transport/messages/BatchMessage.java @@ -20,13 +20,13 @@ package org.apache.cassandra.transport.messages; import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.List; -import java.util.UUID; import com.google.common.collect.ImmutableMap; import io.netty.buffer.ByteBuf; import org.apache.cassandra.audit.AuditLogEntry; import org.apache.cassandra.audit.AuditLogEntryType; +import org.apache.cassandra.audit.AuditLogManager; import org.apache.cassandra.cql3.Attributes; import org.apache.cassandra.cql3.BatchQueryOptions; import org.apache.cassandra.cql3.CQLStatement; @@ -47,7 +47,6 @@ import org.apache.cassandra.transport.ProtocolException; import org.apache.cassandra.transport.ProtocolVersion; import org.apache.cassandra.utils.JVMStabilityInspector; import org.apache.cassandra.utils.MD5Digest; -import org.apache.cassandra.utils.UUIDGen; public class BatchMessage extends Message.Request { @@ -157,30 +156,21 @@ public class BatchMessage extends Message.Request this.options = options; } - public Message.Response execute(QueryState state, long queryStartNanoTime) + @Override + protected boolean isTraceable() { - try - { - UUID tracingId = null; - if (isTracingRequested()) - { - tracingId = UUIDGen.getTimeUUID(); - state.prepareTracingSession(tracingId); - } - - if (state.traceNextQuery()) - { - state.createTracingSession(getCustomPayload()); + return true; + } - ImmutableMap.Builder<String, String> builder = ImmutableMap.builder(); - if(options.getConsistency() != null) - builder.put("consistency_level", options.getConsistency().name()); - if(options.getSerialConsistency() != null) - builder.put("serial_consistency_level", options.getSerialConsistency().name()); + @Override + protected Message.Response execute(QueryState state, long queryStartNanoTime, boolean traceRequest) + { + AuditLogManager auditLogManager = AuditLogManager.getInstance(); - // TODO we don't have [typed] access to CQL bind variables here. CASSANDRA-4560 is open to add support. - Tracing.instance.begin("Execute batch of CQL3 queries", state.getClientAddress(), builder.build()); - } + try + { + if (traceRequest) + traceQuery(state); QueryHandler handler = ClientState.getCQLQueryHandler(); List<QueryHandler.Prepared> prepared = new ArrayList<>(queryOrIdList.size()); @@ -227,39 +217,44 @@ public class BatchMessage extends Message.Request // (and no value would be really correct, so we prefer passing a clearly wrong one). BatchStatement batch = new BatchStatement(batchType, VariableSpecifications.empty(), statements, Attributes.none()); - long fqlTime = isLoggingEnabled ? System.currentTimeMillis() : 0; + long fqlTime = auditLogManager.isLoggingEnabled() ? System.currentTimeMillis() : 0; Message.Response response = handler.processBatch(batch, state, batchOptions, getCustomPayload(), queryStartNanoTime); - if (isLoggingEnabled) - { + if (auditLogManager.isLoggingEnabled()) auditLogManager.logBatch(batchType.name(), queryOrIdList, values, prepared, options, state, fqlTime); - } - - - if (tracingId != null) - response.setTracingId(tracingId); return response; } catch (Exception e) { - if (auditLogEnabled) - { - AuditLogEntry entry = new AuditLogEntry.Builder(state.getClientState()) - .setOperation(getAuditString()) - .setOptions(options) - .setType(AuditLogEntryType.BATCH) - .build(); - auditLogManager.log(entry, e); - } - + if (auditLogManager.isAuditingEnabled()) + logException(state, e); JVMStabilityInspector.inspectThrowable(e); return ErrorMessage.fromException(e); } - finally - { - Tracing.instance.stopSession(); - } + } + + private void traceQuery(QueryState state) + { + ImmutableMap.Builder<String, String> builder = ImmutableMap.builder(); + if (options.getConsistency() != null) + builder.put("consistency_level", options.getConsistency().name()); + if (options.getSerialConsistency() != null) + builder.put("serial_consistency_level", options.getSerialConsistency().name()); + + // TODO we don't have [typed] access to CQL bind variables here. CASSANDRA-4560 is open to add support. + Tracing.instance.begin("Execute batch of CQL3 queries", state.getClientAddress(), builder.build()); + } + + private void logException(QueryState state, Exception e) + { + AuditLogEntry entry = + new AuditLogEntry.Builder(state) + .setOperation(getAuditString()) + .setOptions(options) + .setType(AuditLogEntryType.BATCH) + .build(); + AuditLogManager.getInstance().log(entry, e); } @Override @@ -278,10 +273,6 @@ public class BatchMessage extends Message.Request private String getAuditString() { - StringBuilder sb = new StringBuilder(); - sb.append("BATCH of ["); - sb.append(queryOrIdList.size()); - sb.append("] statements at consistency ").append(options.getConsistency()); - return sb.toString(); + return String.format("BATCH of %d statements at consistency %s", queryOrIdList.size(), options.getConsistency()); } } http://git-wip-us.apache.org/repos/asf/cassandra/blob/a09dc3a5/src/java/org/apache/cassandra/transport/messages/ExecuteMessage.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/transport/messages/ExecuteMessage.java b/src/java/org/apache/cassandra/transport/messages/ExecuteMessage.java index 6c0b77a..ff888f2 100644 --- a/src/java/org/apache/cassandra/transport/messages/ExecuteMessage.java +++ b/src/java/org/apache/cassandra/transport/messages/ExecuteMessage.java @@ -17,8 +17,6 @@ */ package org.apache.cassandra.transport.messages; -import java.util.UUID; - import com.google.common.collect.ImmutableMap; import io.netty.buffer.ByteBuf; @@ -28,7 +26,6 @@ import org.apache.cassandra.cql3.CQLStatement; import org.apache.cassandra.cql3.ColumnSpecification; import org.apache.cassandra.cql3.QueryHandler; import org.apache.cassandra.cql3.QueryOptions; -import org.apache.cassandra.cql3.QueryProcessor; import org.apache.cassandra.cql3.ResultSet; import org.apache.cassandra.exceptions.PreparedQueryNotFoundException; import org.apache.cassandra.service.ClientState; @@ -40,7 +37,6 @@ import org.apache.cassandra.transport.ProtocolException; import org.apache.cassandra.transport.ProtocolVersion; import org.apache.cassandra.utils.JVMStabilityInspector; import org.apache.cassandra.utils.MD5Digest; -import org.apache.cassandra.utils.UUIDGen; public class ExecuteMessage extends Message.Request { @@ -108,12 +104,21 @@ public class ExecuteMessage extends Message.Request this.resultMetadataId = resultMetadataId; } - public Message.Response execute(QueryState state, long queryStartNanoTime) + @Override + protected boolean isTraceable() { + return true; + } + + @Override + protected Message.Response execute(QueryState state, long queryStartNanoTime, boolean traceRequest) + { + AuditLogManager auditLogManager = AuditLogManager.getInstance(); + try { QueryHandler handler = ClientState.getCQLQueryHandler(); - QueryProcessor.Prepared prepared = handler.getPrepared(statementId); + QueryHandler.Prepared prepared = handler.getPrepared(statementId); if (prepared == null) throw new PreparedQueryNotFoundException(statementId); @@ -123,63 +128,19 @@ public class ExecuteMessage extends Message.Request if (options.getPageSize() == 0) throw new ProtocolException("The page size cannot be 0"); - UUID tracingId = null; - if (isTracingRequested()) - { - tracingId = UUIDGen.getTimeUUID(); - state.prepareTracingSession(tracingId); - } - - if (state.traceNextQuery()) - { - state.createTracingSession(getCustomPayload()); - - ImmutableMap.Builder<String, String> builder = ImmutableMap.builder(); - if (options.getPageSize() > 0) - builder.put("page_size", Integer.toString(options.getPageSize())); - if(options.getConsistency() != null) - builder.put("consistency_level", options.getConsistency().name()); - if(options.getSerialConsistency() != null) - builder.put("serial_consistency_level", options.getSerialConsistency().name()); - builder.put("query", prepared.rawCQLStatement); - - for(int i = 0; i < statement.getBindVariables().size(); i++) - { - ColumnSpecification cs = statement.getBindVariables().get(i); - String boundName = cs.name.toString(); - String boundValue = cs.type.asCQL3Type().toCQLLiteral(options.getValues().get(i), options.getProtocolVersion()); - if ( boundValue.length() > 1000 ) - { - boundValue = boundValue.substring(0, 1000) + "...'"; - } - - //Here we prefix boundName with the index to avoid possible collission in builder keys due to - //having multiple boundValues for the same variable - builder.put("bound_var_" + Integer.toString(i) + "_" + boundName, boundValue); - } - - Tracing.instance.begin("Execute CQL3 prepared query", state.getClientAddress(), builder.build()); - } + if (traceRequest) + traceQuery(state, prepared); // Some custom QueryHandlers are interested by the bound names. We provide them this information // by wrapping the QueryOptions. QueryOptions queryOptions = QueryOptions.addColumnSpecifications(options, prepared.statement.getBindVariables()); - long fqlTime = isLoggingEnabled ? System.currentTimeMillis() : 0; + long requestStartTime = auditLogManager.isLoggingEnabled() ? System.currentTimeMillis() : 0L; + Message.Response response = handler.processPrepared(statement, state, queryOptions, getCustomPayload(), queryStartNanoTime); - if (isLoggingEnabled) - { - AuditLogEntry auditEntry = new AuditLogEntry.Builder(state.getClientState()) - .setType(statement.getAuditLogContext().auditLogEntryType) - .setOperation(prepared.rawCQLStatement) - .setTimestamp(fqlTime) - .setScope(statement) - .setKeyspace(state, statement) - .setOptions(options) - .build(); - AuditLogManager.getInstance().log(auditEntry); - } + if (auditLogManager.isLoggingEnabled()) + logSuccess(state, prepared, requestStartTime); if (response instanceof ResultMessage.Rows) { @@ -211,52 +172,90 @@ public class ExecuteMessage extends Message.Request } } - if (tracingId != null) - response.setTracingId(tracingId); - return response; } catch (Exception e) { - if (auditLogEnabled) - { - if (e instanceof PreparedQueryNotFoundException) - { - AuditLogEntry auditLogEntry = new AuditLogEntry.Builder(state.getClientState()) - .setOperation(toString()) - .setOptions(options) - .build(); - auditLogManager.log(auditLogEntry, e); - } - else - { - QueryHandler.Prepared prepared = ClientState.getCQLQueryHandler().getPrepared(statementId); - if (prepared != null) - { - AuditLogEntry auditLogEntry = new AuditLogEntry.Builder(state.getClientState()) - .setOperation(toString()) - .setType(prepared.statement.getAuditLogContext().auditLogEntryType) - .setScope(prepared.statement) - .setKeyspace(state, prepared.statement) - .setOptions(options) - .build(); - auditLogManager.log(auditLogEntry, e); - } - } - } - + if (auditLogManager.isAuditingEnabled()) + logException(state, e); JVMStabilityInspector.inspectThrowable(e); return ErrorMessage.fromException(e); } - finally + } + + private void traceQuery(QueryState state, QueryHandler.Prepared prepared) + { + ImmutableMap.Builder<String, String> builder = ImmutableMap.builder(); + if (options.getPageSize() > 0) + builder.put("page_size", Integer.toString(options.getPageSize())); + if (options.getConsistency() != null) + builder.put("consistency_level", options.getConsistency().name()); + if (options.getSerialConsistency() != null) + builder.put("serial_consistency_level", options.getSerialConsistency().name()); + + builder.put("query", prepared.rawCQLStatement); + + for (int i = 0; i < prepared.statement.getBindVariables().size(); i++) + { + ColumnSpecification cs = prepared.statement.getBindVariables().get(i); + String boundName = cs.name.toString(); + String boundValue = cs.type.asCQL3Type().toCQLLiteral(options.getValues().get(i), options.getProtocolVersion()); + if (boundValue.length() > 1000) + boundValue = boundValue.substring(0, 1000) + "...'"; + + //Here we prefix boundName with the index to avoid possible collission in builder keys due to + //having multiple boundValues for the same variable + builder.put("bound_var_" + i + '_' + boundName, boundValue); + } + + Tracing.instance.begin("Execute CQL3 prepared query", state.getClientAddress(), builder.build()); + } + + private void logSuccess(QueryState state, QueryHandler.Prepared prepared, long requestStartTime) + { + AuditLogEntry entry = + new AuditLogEntry.Builder(state) + .setType(prepared.statement.getAuditLogContext().auditLogEntryType) + .setOperation(prepared.rawCQLStatement) + .setTimestamp(requestStartTime) + .setScope(prepared.statement) + .setKeyspace(state, prepared.statement) + .setOptions(options) + .build(); + AuditLogManager.getInstance().log(entry); + } + + private void logException(QueryState state, Exception e) + { + if (e instanceof PreparedQueryNotFoundException) + { + AuditLogEntry entry = + new AuditLogEntry.Builder(state) + .setOperation(toString()) + .setOptions(options) + .build(); + AuditLogManager.getInstance().log(entry, e); + return; + } + + QueryHandler.Prepared prepared = ClientState.getCQLQueryHandler().getPrepared(statementId); + if (prepared != null) { - Tracing.instance.stopSession(); + AuditLogEntry entry = + new AuditLogEntry.Builder(state) + .setOperation(toString()) + .setType(prepared.statement.getAuditLogContext().auditLogEntryType) + .setScope(prepared.statement) + .setKeyspace(state, prepared.statement) + .setOptions(options) + .build(); + AuditLogManager.getInstance().log(entry, e); } } @Override public String toString() { - return "EXECUTE " + statementId + " with " + options.getValues().size() + " values at consistency " + options.getConsistency(); + return String.format("EXECUTE %s with %d values at consistency %s", statementId, options.getValues().size(), options.getConsistency()); } } http://git-wip-us.apache.org/repos/asf/cassandra/blob/a09dc3a5/src/java/org/apache/cassandra/transport/messages/OptionsMessage.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/transport/messages/OptionsMessage.java b/src/java/org/apache/cassandra/transport/messages/OptionsMessage.java index 914ccb1..2b8e695 100644 --- a/src/java/org/apache/cassandra/transport/messages/OptionsMessage.java +++ b/src/java/org/apache/cassandra/transport/messages/OptionsMessage.java @@ -57,7 +57,8 @@ public class OptionsMessage extends Message.Request super(Message.Type.OPTIONS); } - public Message.Response execute(QueryState state, long queryStartNanoTime) + @Override + protected Message.Response execute(QueryState state, long queryStartNanoTime, boolean traceRequest) { List<String> cqlVersions = new ArrayList<String>(); cqlVersions.add(QueryProcessor.CQL_VERSION.toString()); http://git-wip-us.apache.org/repos/asf/cassandra/blob/a09dc3a5/src/java/org/apache/cassandra/transport/messages/PrepareMessage.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/transport/messages/PrepareMessage.java b/src/java/org/apache/cassandra/transport/messages/PrepareMessage.java index 4ab6e0b..d9d3ed8 100644 --- a/src/java/org/apache/cassandra/transport/messages/PrepareMessage.java +++ b/src/java/org/apache/cassandra/transport/messages/PrepareMessage.java @@ -17,13 +17,12 @@ */ package org.apache.cassandra.transport.messages; -import java.util.UUID; - import com.google.common.collect.ImmutableMap; import io.netty.buffer.ByteBuf; import org.apache.cassandra.audit.AuditLogEntry; import org.apache.cassandra.audit.AuditLogEntryType; +import org.apache.cassandra.audit.AuditLogManager; import org.apache.cassandra.cql3.CQLStatement; import org.apache.cassandra.cql3.QueryProcessor; import org.apache.cassandra.service.ClientState; @@ -33,7 +32,6 @@ import org.apache.cassandra.transport.CBUtil; import org.apache.cassandra.transport.Message; import org.apache.cassandra.transport.ProtocolVersion; import org.apache.cassandra.utils.JVMStabilityInspector; -import org.apache.cassandra.utils.UUIDGen; public class PrepareMessage extends Message.Request { @@ -97,61 +95,62 @@ public class PrepareMessage extends Message.Request this.keyspace = keyspace; } - public Message.Response execute(QueryState state, long queryStartNanoTime) + @Override + protected boolean isTraceable() + { + return true; + } + + @Override + protected Message.Response execute(QueryState state, long queryStartNanoTime, boolean traceRequest) { + AuditLogManager auditLogManager = AuditLogManager.getInstance(); + try { - UUID tracingId = null; - if (isTracingRequested()) - { - tracingId = UUIDGen.getTimeUUID(); - state.prepareTracingSession(tracingId); - } - - if (state.traceNextQuery()) - { - state.createTracingSession(getCustomPayload()); + if (traceRequest) Tracing.instance.begin("Preparing CQL3 query", state.getClientAddress(), ImmutableMap.of("query", query)); - } - Message.Response response = ClientState.getCQLQueryHandler().prepare(query, - state.getClientState().cloneWithKeyspaceIfSet(keyspace), - getCustomPayload()); - if (auditLogEnabled) - { - CQLStatement parsedStmt = QueryProcessor.parseStatement(query, state.getClientState()); - AuditLogEntry auditLogEntry = new AuditLogEntry.Builder(state.getClientState()) - .setOperation(query) - .setType(AuditLogEntryType.PREPARE_STATEMENT) - .setScope(parsedStmt) - .setKeyspace(parsedStmt) - .build(); - auditLogManager.log(auditLogEntry); - } + ClientState clientState = state.getClientState().cloneWithKeyspaceIfSet(keyspace); + Message.Response response = ClientState.getCQLQueryHandler().prepare(query, clientState, getCustomPayload()); - if (tracingId != null) - response.setTracingId(tracingId); + if (auditLogManager.isAuditingEnabled()) + logSuccess(state); return response; } catch (Exception e) { - if (auditLogEnabled) - { - AuditLogEntry auditLogEntry = new AuditLogEntry.Builder(state.getClientState()) - .setOperation(query) - .setKeyspace(keyspace) - .setType(AuditLogEntryType.PREPARE_STATEMENT) - .build(); - auditLogManager.log(auditLogEntry, e); - } + if (auditLogManager.isAuditingEnabled()) + logException(state, e); JVMStabilityInspector.inspectThrowable(e); return ErrorMessage.fromException(e); } - finally - { - Tracing.instance.stopSession(); - } + } + + private void logSuccess(QueryState state) + { + // The statement gets parsed twice. Not a big deal given that this is PREPARE, but still, why? + CQLStatement statement = QueryProcessor.parseStatement(query, state.getClientState()); + AuditLogEntry entry = + new AuditLogEntry.Builder(state) + .setOperation(query) + .setType(AuditLogEntryType.PREPARE_STATEMENT) + .setScope(statement) + .setKeyspace(statement) + .build(); + AuditLogManager.getInstance().log(entry); + } + + private void logException(QueryState state, Exception e) + { + AuditLogEntry entry = + new AuditLogEntry.Builder(state) + .setOperation(query) + .setKeyspace(keyspace) + .setType(AuditLogEntryType.PREPARE_STATEMENT) + .build(); + AuditLogManager.getInstance().log(entry, e); } @Override http://git-wip-us.apache.org/repos/asf/cassandra/blob/a09dc3a5/src/java/org/apache/cassandra/transport/messages/QueryMessage.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/transport/messages/QueryMessage.java b/src/java/org/apache/cassandra/transport/messages/QueryMessage.java index 4f42b85..7b80200 100644 --- a/src/java/org/apache/cassandra/transport/messages/QueryMessage.java +++ b/src/java/org/apache/cassandra/transport/messages/QueryMessage.java @@ -17,8 +17,6 @@ */ package org.apache.cassandra.transport.messages; -import java.util.UUID; - import com.google.common.collect.ImmutableMap; import io.netty.buffer.ByteBuf; @@ -37,7 +35,6 @@ import org.apache.cassandra.transport.Message; import org.apache.cassandra.transport.ProtocolException; import org.apache.cassandra.transport.ProtocolVersion; import org.apache.cassandra.utils.JVMStabilityInspector; -import org.apache.cassandra.utils.UUIDGen; /** * A CQL query @@ -87,86 +84,91 @@ public class QueryMessage extends Message.Request this.options = options; } - public Message.Response execute(QueryState state, long queryStartNanoTime) + @Override + protected boolean isTraceable() { + return true; + } + + @Override + protected Message.Response execute(QueryState state, long queryStartNanoTime, boolean traceRequest) + { + AuditLogManager auditLogManager = AuditLogManager.getInstance(); + try { if (options.getPageSize() == 0) throw new ProtocolException("The page size cannot be 0"); - UUID tracingId = null; - if (isTracingRequested()) - { - tracingId = UUIDGen.getTimeUUID(); - state.prepareTracingSession(tracingId); - } + if (traceRequest) + traceQuery(state); - if (state.traceNextQuery()) - { - state.createTracingSession(getCustomPayload()); - - ImmutableMap.Builder<String, String> builder = ImmutableMap.builder(); - builder.put("query", query); - if (options.getPageSize() > 0) - builder.put("page_size", Integer.toString(options.getPageSize())); - if(options.getConsistency() != null) - builder.put("consistency_level", options.getConsistency().name()); - if(options.getSerialConsistency() != null) - builder.put("serial_consistency_level", options.getSerialConsistency().name()); - - Tracing.instance.begin("Execute CQL3 query", state.getClientAddress(), builder.build()); - } + long queryStartTime = auditLogManager.isLoggingEnabled() ? System.currentTimeMillis() : 0L; - long fqlTime = isLoggingEnabled ? System.currentTimeMillis() : 0; Message.Response response = ClientState.getCQLQueryHandler().process(query, state, options, getCustomPayload(), queryStartNanoTime); - if (isLoggingEnabled) - { - CQLStatement parsedStatement = QueryProcessor.parseStatement(query, state.getClientState()); - AuditLogEntry auditEntry = new AuditLogEntry.Builder(state.getClientState()) - .setType(parsedStatement.getAuditLogContext().auditLogEntryType) - .setOperation(query) - .setTimestamp(fqlTime) - .setScope(parsedStatement) - .setKeyspace(state, parsedStatement) - .setOptions(options) - .build(); - AuditLogManager.getInstance().log(auditEntry); - - } + if (auditLogManager.isLoggingEnabled()) + logSuccess(state, queryStartTime); if (options.skipMetadata() && response instanceof ResultMessage.Rows) ((ResultMessage.Rows)response).result.metadata.setSkipMetadata(); - if (tracingId != null) - response.setTracingId(tracingId); - return response; } catch (Exception e) { - if (auditLogEnabled) - { - AuditLogEntry auditLogEntry = new AuditLogEntry.Builder(state.getClientState()) - .setOperation(query) - .setOptions(options) - .build(); - auditLogManager.log(auditLogEntry, e); - } + if (auditLogManager.isLoggingEnabled()) + logException(state, e); JVMStabilityInspector.inspectThrowable(e); if (!((e instanceof RequestValidationException) || (e instanceof RequestExecutionException))) logger.error("Unexpected error during query", e); return ErrorMessage.fromException(e); } - finally - { - Tracing.instance.stopSession(); - } + } + + private void traceQuery(QueryState state) + { + ImmutableMap.Builder<String, String> builder = ImmutableMap.builder(); + builder.put("query", query); + if (options.getPageSize() > 0) + builder.put("page_size", Integer.toString(options.getPageSize())); + if (options.getConsistency() != null) + builder.put("consistency_level", options.getConsistency().name()); + if (options.getSerialConsistency() != null) + builder.put("serial_consistency_level", options.getSerialConsistency().name()); + + Tracing.instance.begin("Execute CQL3 query", state.getClientAddress(), builder.build()); + } + + private void logSuccess(QueryState state, long queryStartTime) + { + // FIXME: we are parsing the statement twice if audit logging is enabled. Why? + CQLStatement statement = QueryProcessor.parseStatement(query, state.getClientState()); + AuditLogEntry entry = + new AuditLogEntry.Builder(state) + .setType(statement.getAuditLogContext().auditLogEntryType) + .setOperation(query) + .setTimestamp(queryStartTime) + .setScope(statement) + .setKeyspace(state, statement) + .setOptions(options) + .build(); + AuditLogManager.getInstance().log(entry); + } + + private void logException(QueryState state, Exception e) + { + AuditLogEntry entry = + new AuditLogEntry.Builder(state) + .setOperation(query) + .setOptions(options) + .build(); + AuditLogManager.getInstance().log(entry, e); } @Override public String toString() { - return "QUERY " + query + "[pageSize = " + options.getPageSize() + "]"; + return String.format("QUERY %s [pageSize = %d]", query, options.getPageSize()); } } http://git-wip-us.apache.org/repos/asf/cassandra/blob/a09dc3a5/src/java/org/apache/cassandra/transport/messages/RegisterMessage.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/transport/messages/RegisterMessage.java b/src/java/org/apache/cassandra/transport/messages/RegisterMessage.java index 2356dae..f8eb55e 100644 --- a/src/java/org/apache/cassandra/transport/messages/RegisterMessage.java +++ b/src/java/org/apache/cassandra/transport/messages/RegisterMessage.java @@ -62,7 +62,8 @@ public class RegisterMessage extends Message.Request this.eventTypes = eventTypes; } - public Response execute(QueryState state, long queryStartNanoTime) + @Override + protected Response execute(QueryState state, long queryStartNanoTime, boolean traceRequest) { assert connection instanceof ServerConnection; Connection.Tracker tracker = connection.getTracker(); http://git-wip-us.apache.org/repos/asf/cassandra/blob/a09dc3a5/src/java/org/apache/cassandra/transport/messages/StartupMessage.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/transport/messages/StartupMessage.java b/src/java/org/apache/cassandra/transport/messages/StartupMessage.java index f72121a..01b9331 100644 --- a/src/java/org/apache/cassandra/transport/messages/StartupMessage.java +++ b/src/java/org/apache/cassandra/transport/messages/StartupMessage.java @@ -66,7 +66,8 @@ public class StartupMessage extends Message.Request this.options = options; } - public Message.Response execute(QueryState state, long queryStartNanoTime) + @Override + protected Message.Response execute(QueryState state, long queryStartNanoTime, boolean traceRequest) { String cqlVersion = options.get(CQL_VERSION); if (cqlVersion == null) --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
