Repository: cassandra Updated Branches: refs/heads/trunk 2d8ae8373 -> 53c71949d
Clean up and refactor client metrics patch by Aleksey Yeschenko and Chris Lohfink for CASSANDRA-14524 Co-authored-by: Chris Lohfink <[email protected]> Project: http://git-wip-us.apache.org/repos/asf/cassandra/repo Commit: http://git-wip-us.apache.org/repos/asf/cassandra/commit/53c71949 Tree: http://git-wip-us.apache.org/repos/asf/cassandra/tree/53c71949 Diff: http://git-wip-us.apache.org/repos/asf/cassandra/diff/53c71949 Branch: refs/heads/trunk Commit: 53c71949d4a49d6062e43ff3b1a9ca3e94496cfb Parents: 2d8ae83 Author: Aleksey Yeschenko <[email protected]> Authored: Sat Jun 16 22:05:51 2018 -0500 Committer: Aleksey Yeshchenko <[email protected]> Committed: Thu Jun 21 16:01:18 2018 +0100 ---------------------------------------------------------------------- CHANGES.txt | 1 + .../apache/cassandra/metrics/AuthMetrics.java | 40 ------ .../apache/cassandra/metrics/ClientMetrics.java | 109 +++++++++++--- .../service/NativeTransportService.java | 53 +------ .../org/apache/cassandra/tools/NodeProbe.java | 2 +- .../cassandra/tools/nodetool/ClientStats.java | 23 ++- .../apache/cassandra/transport/ClientStat.java | 56 ++++++++ .../cassandra/transport/ConnectedClient.java | 144 +++++++++++++++++++ .../cassandra/transport/ConnectionStage.java | 23 +++ .../transport/ProtocolVersionTracker.java | 75 ++++------ .../org/apache/cassandra/transport/Server.java | 71 +++------ .../cassandra/transport/ServerConnection.java | 31 ++-- .../transport/messages/AuthResponse.java | 6 +- .../transport/ProtocolVersionTrackerTest.java | 23 ++- 14 files changed, 409 insertions(+), 248 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/cassandra/blob/53c71949/CHANGES.txt ---------------------------------------------------------------------- diff --git a/CHANGES.txt b/CHANGES.txt index 041a655..598eaff 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 4.0 + * Clean up and refactor client metrics (CASSANDRA-14524) * Nodetool import row cache invalidation races with adding sstables to tracker (CASSANDRA-14529) * Fix assertions in LWTs after TableMetadata was made immutable (CASSANDRA-14356) * Abort compactions quicker (CASSANDRA-14397) http://git-wip-us.apache.org/repos/asf/cassandra/blob/53c71949/src/java/org/apache/cassandra/metrics/AuthMetrics.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/metrics/AuthMetrics.java b/src/java/org/apache/cassandra/metrics/AuthMetrics.java deleted file mode 100644 index 126738c..0000000 --- a/src/java/org/apache/cassandra/metrics/AuthMetrics.java +++ /dev/null @@ -1,40 +0,0 @@ -package org.apache.cassandra.metrics; - -import com.codahale.metrics.Meter; - -/** - * Metrics about authentication - */ -public class AuthMetrics -{ - - public static final AuthMetrics instance = new AuthMetrics(); - - public static void init() - { - // no-op, just used to force instance creation - } - - /** Number and rate of successful logins */ - protected final Meter success; - - /** Number and rate of login failures */ - protected final Meter failure; - - private AuthMetrics() - { - - success = ClientMetrics.instance.addMeter("AuthSuccess"); - failure = ClientMetrics.instance.addMeter("AuthFailure"); - } - - public void markSuccess() - { - success.mark(); - } - - public void markFailure() - { - failure.mark(); - } -} http://git-wip-us.apache.org/repos/asf/cassandra/blob/53c71949/src/java/org/apache/cassandra/metrics/ClientMetrics.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/metrics/ClientMetrics.java b/src/java/org/apache/cassandra/metrics/ClientMetrics.java index 8ca3480..5e7720a 100644 --- a/src/java/org/apache/cassandra/metrics/ClientMetrics.java +++ b/src/java/org/apache/cassandra/metrics/ClientMetrics.java @@ -18,38 +18,113 @@ */ package org.apache.cassandra.metrics; -import java.util.concurrent.Callable; +import java.util.*; import com.codahale.metrics.Gauge; import com.codahale.metrics.Meter; +import org.apache.cassandra.transport.ClientStat; +import org.apache.cassandra.transport.ConnectedClient; +import org.apache.cassandra.transport.Server; import static org.apache.cassandra.metrics.CassandraMetricsRegistry.Metrics; - -public class ClientMetrics +public final class ClientMetrics { - private static final MetricNameFactory factory = new DefaultNameFactory("Client"); - public static final ClientMetrics instance = new ClientMetrics(); - + + private static final MetricNameFactory factory = new DefaultNameFactory("Client"); + + private volatile boolean initialized = false; + private Collection<Server> servers = Collections.emptyList(); + + private Meter authSuccess; + private Meter authFailure; + private ClientMetrics() { } - public <T> void addGauge(String name, final Callable<T> provider) + public void markAuthSuccess() + { + authSuccess.mark(); + } + + public void markAuthFailure() + { + authFailure.mark(); + } + + public synchronized void init(Collection<Server> servers) + { + if (initialized) + return; + + this.servers = servers; + + registerGauge("connectedNativeClients", this::countConnectedClients); + registerGauge("connectedNativeClientsByUser", this::countConnectedClientsByUser); + registerGauge("connections", this::connectedClients); + registerGauge("clientsByProtocolVersion", this::recentClientStats); + + authSuccess = registerMeter("AuthSuccess"); + authFailure = registerMeter("AuthFailure"); + + initialized = true; + } + + private int countConnectedClients() + { + int count = 0; + + for (Server server : servers) + count += server.countConnectedClients(); + + return count; + } + + private Map<String, Integer> countConnectedClientsByUser() + { + Map<String, Integer> counts = new HashMap<>(); + + for (Server server : servers) + { + server.countConnectedClientsByUser() + .forEach((username, count) -> counts.put(username, counts.getOrDefault(username, 0) + count)); + } + + return counts; + } + + private List<Map<String, String>> connectedClients() + { + List<Map<String, String>> clients = new ArrayList<>(); + + for (Server server : servers) + for (ConnectedClient client : server.getConnectedClients()) + clients.add(client.asMap()); + + return clients; + } + + private List<Map<String, String>> recentClientStats() + { + List<Map<String, String>> stats = new ArrayList<>(); + + for (Server server : servers) + for (ClientStat stat : server.recentClientStats()) + stats.add(stat.asMap()); + + stats.sort(Comparator.comparing(map -> map.get(ClientStat.PROTOCOL_VERSION))); + + return stats; + } + + private <T> Gauge<T> registerGauge(String name, Gauge<T> gauge) { - Metrics.register(factory.createMetricName(name), (Gauge<T>) () -> { - try - { - return provider.call(); - } catch (Exception e) - { - throw new RuntimeException(e); - } - }); + return Metrics.register(factory.createMetricName(name), gauge); } - public Meter addMeter(String name) + private Meter registerMeter(String name) { return Metrics.meter(factory.createMetricName(name)); } http://git-wip-us.apache.org/repos/asf/cassandra/blob/53c71949/src/java/org/apache/cassandra/service/NativeTransportService.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/service/NativeTransportService.java b/src/java/org/apache/cassandra/service/NativeTransportService.java index 39b334e..79acab1 100644 --- a/src/java/org/apache/cassandra/service/NativeTransportService.java +++ b/src/java/org/apache/cassandra/service/NativeTransportService.java @@ -18,14 +18,9 @@ package org.apache.cassandra.service; import java.net.InetAddress; -import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; import java.util.Collections; -import java.util.HashMap; -import java.util.List; -import java.util.Map; -import java.util.Map.Entry; import java.util.concurrent.TimeUnit; import com.google.common.annotations.VisibleForTesting; @@ -39,7 +34,6 @@ import io.netty.channel.epoll.EpollEventLoopGroup; import io.netty.channel.nio.NioEventLoopGroup; import io.netty.util.concurrent.EventExecutor; import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.metrics.AuthMetrics; import org.apache.cassandra.metrics.ClientMetrics; import org.apache.cassandra.transport.RequestThreadPoolExecutor; import org.apache.cassandra.transport.Server; @@ -114,52 +108,7 @@ public class NativeTransportService } } - // register metrics - ClientMetrics.instance.addGauge("connectedNativeClients", () -> - { - int ret = 0; - for (Server server : servers) - ret += server.getConnectedClients(); - return ret; - }); - ClientMetrics.instance.addGauge("connectedNativeClientsByUser", () -> - { - Map<String, Integer> result = new HashMap<>(); - for (Server server : servers) - { - for (Entry<String, Integer> e : server.getConnectedClientsByUser().entrySet()) - { - String user = e.getKey(); - result.put(user, result.getOrDefault(user, 0) + e.getValue()); - } - } - return result; - }); - - ClientMetrics.instance.addGauge("connections", () -> - { - List<Map<String, String>> result = new ArrayList<>(); - for (Server server : servers) - { - for (Map<String, String> e : server.getConnectionStates()) - { - result.add(e); - } - } - return result; - }); - - ClientMetrics.instance.addGauge("clientsByProtocolVersion", () -> - { - List<Map<String, String>> result = new ArrayList<>(); - for (Server server : servers) - { - result.addAll(server.getClientsByProtocolVersion()); - } - return result; - }); - - AuthMetrics.init(); + ClientMetrics.instance.init(servers); initialized = true; } http://git-wip-us.apache.org/repos/asf/cassandra/blob/53c71949/src/java/org/apache/cassandra/tools/NodeProbe.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/tools/NodeProbe.java b/src/java/org/apache/cassandra/tools/NodeProbe.java index 01769b4..caaa337 100644 --- a/src/java/org/apache/cassandra/tools/NodeProbe.java +++ b/src/java/org/apache/cassandra/tools/NodeProbe.java @@ -1517,7 +1517,7 @@ public class NodeProbe implements AutoCloseable /** * Retrieve Proxy metrics - * @param connections, connectedNativeClients, connectedNativeClientsByUser + * @param connections, connectedNativeClients, connectedNativeClientsByUser, clientsByProtocolVersion */ public Object getClientMetric(String metricName) { http://git-wip-us.apache.org/repos/asf/cassandra/blob/53c71949/src/java/org/apache/cassandra/tools/nodetool/ClientStats.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/tools/nodetool/ClientStats.java b/src/java/org/apache/cassandra/tools/nodetool/ClientStats.java index 5bd5da1..3bf46b4 100644 --- a/src/java/org/apache/cassandra/tools/nodetool/ClientStats.java +++ b/src/java/org/apache/cassandra/tools/nodetool/ClientStats.java @@ -23,12 +23,13 @@ import java.util.List; import java.util.Map; import java.util.Map.Entry; +import io.airlift.airline.Command; +import io.airlift.airline.Option; import org.apache.cassandra.tools.NodeProbe; import org.apache.cassandra.tools.NodeTool.NodeToolCmd; import org.apache.cassandra.tools.nodetool.formatter.TableBuilder; - -import io.airlift.airline.Command; -import io.airlift.airline.Option; +import org.apache.cassandra.transport.ClientStat; +import org.apache.cassandra.transport.ConnectedClient; @Command(name = "clientstats", description = "Print information about connected clients") public class ClientStats extends NodeToolCmd @@ -68,7 +69,9 @@ public class ClientStats extends NodeToolCmd for (Map<String, String> client : clients) { - table.add(client.get("protocolVersion"), client.get("inetAddress"), sdf.format(new Date(Long.valueOf(client.get("lastSeenTime"))))); + table.add(client.get(ClientStat.PROTOCOL_VERSION), + client.get(ClientStat.INET_ADDRESS), + sdf.format(new Date(Long.valueOf(client.get(ClientStat.LAST_SEEN_TIME))))); } table.printTo(System.out); @@ -87,8 +90,16 @@ public class ClientStats extends NodeToolCmd table.add("Address", "SSL", "Cipher", "Protocol", "Version", "User", "Keyspace", "Requests", "Driver-Name", "Driver-Version"); for (Map<String, String> conn : clients) { - table.add(conn.get("address"), conn.get("ssl"), conn.get("cipher"), conn.get("protocol"), conn.get("version"), - conn.get("user"), conn.get("keyspace"), conn.get("requests"), conn.get("driverName"), conn.get("driverVersion")); + table.add(conn.get(ConnectedClient.ADDRESS), + conn.get(ConnectedClient.SSL), + conn.get(ConnectedClient.CIPHER), + conn.get(ConnectedClient.PROTOCOL), + conn.get(ConnectedClient.VERSION), + conn.get(ConnectedClient.USER), + conn.get(ConnectedClient.KEYSPACE), + conn.get(ConnectedClient.REQUESTS), + conn.get(ConnectedClient.DRIVER_NAME), + conn.get(ConnectedClient.DRIVER_VERSION)); } table.printTo(System.out); System.out.println(); http://git-wip-us.apache.org/repos/asf/cassandra/blob/53c71949/src/java/org/apache/cassandra/transport/ClientStat.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/transport/ClientStat.java b/src/java/org/apache/cassandra/transport/ClientStat.java new file mode 100644 index 0000000..7e78597 --- /dev/null +++ b/src/java/org/apache/cassandra/transport/ClientStat.java @@ -0,0 +1,56 @@ +/* + * 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.transport; + +import java.net.InetAddress; +import java.util.Map; + +import com.google.common.collect.ImmutableMap; + +public final class ClientStat +{ + public static final String INET_ADDRESS = "inetAddress"; + public static final String PROTOCOL_VERSION = "protocolVersion"; + public static final String LAST_SEEN_TIME = "lastSeenTime"; + + final InetAddress remoteAddress; + final ProtocolVersion protocolVersion; + final long lastSeenTime; + + ClientStat(InetAddress remoteAddress, ProtocolVersion protocolVersion, long lastSeenTime) + { + this.remoteAddress = remoteAddress; + this.lastSeenTime = lastSeenTime; + this.protocolVersion = protocolVersion; + } + + @Override + public String toString() + { + return String.format("ClientStat{%s, %s, %d}", remoteAddress, protocolVersion, lastSeenTime); + } + + public Map<String, String> asMap() + { + return ImmutableMap.<String, String>builder() + .put(INET_ADDRESS, remoteAddress.toString()) + .put(PROTOCOL_VERSION, protocolVersion.toString()) + .put(LAST_SEEN_TIME, String.valueOf(lastSeenTime)) + .build(); + } +} http://git-wip-us.apache.org/repos/asf/cassandra/blob/53c71949/src/java/org/apache/cassandra/transport/ConnectedClient.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/transport/ConnectedClient.java b/src/java/org/apache/cassandra/transport/ConnectedClient.java new file mode 100644 index 0000000..0776bf8 --- /dev/null +++ b/src/java/org/apache/cassandra/transport/ConnectedClient.java @@ -0,0 +1,144 @@ +/* + * 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.transport; + +import java.net.InetSocketAddress; +import java.util.Map; +import java.util.Optional; + +import com.google.common.collect.ImmutableMap; + +import io.netty.handler.ssl.SslHandler; +import org.apache.cassandra.auth.AuthenticatedUser; +import org.apache.cassandra.service.ClientState; + +public final class ConnectedClient +{ + public static final String ADDRESS = "address"; + public static final String USER = "user"; + public static final String VERSION = "version"; + public static final String DRIVER_NAME = "driverName"; + public static final String DRIVER_VERSION = "driverVersion"; + public static final String REQUESTS = "requests"; + public static final String KEYSPACE = "keyspace"; + public static final String SSL = "ssl"; + public static final String CIPHER = "cipher"; + public static final String PROTOCOL = "protocol"; + + private static final String UNDEFINED = "undefined"; + + private final ServerConnection connection; + + ConnectedClient(ServerConnection connection) + { + this.connection = connection; + } + + public ConnectionStage stage() + { + return connection.stage(); + } + + public InetSocketAddress remoteAddress() + { + return state().getRemoteAddress(); + } + + public Optional<String> username() + { + AuthenticatedUser user = state().getUser(); + + return null != user + ? Optional.of(user.getName()) + : Optional.empty(); + } + + public int protocolVersion() + { + return connection.getVersion().asInt(); + } + + public Optional<String> driverName() + { + return state().getDriverName(); + } + + public Optional<String> driverVersion() + { + return state().getDriverVersion(); + } + + public long requestCount() + { + return connection.requests.getCount(); + } + + public Optional<String> keyspace() + { + return Optional.ofNullable(state().getRawKeyspace()); + } + + public boolean isEncrypted() + { + return null != sslHandler(); + } + + public Optional<String> sslCipherSuite() + { + SslHandler sslHandler = sslHandler(); + + return null != sslHandler + ? Optional.of(sslHandler.engine().getSession().getCipherSuite()) + : Optional.empty(); + } + + public Optional<String> sslProtocol() + { + SslHandler sslHandler = sslHandler(); + + return null != sslHandler + ? Optional.of(sslHandler.engine().getSession().getProtocol()) + : Optional.empty(); + } + + private ClientState state() + { + return connection.getClientState(); + } + + private SslHandler sslHandler() + { + return connection.channel().pipeline().get(SslHandler.class); + } + + public Map<String, String> asMap() + { + return ImmutableMap.<String, String>builder() + .put(ADDRESS, remoteAddress().toString()) + .put(USER, username().orElse(UNDEFINED)) + .put(VERSION, String.valueOf(protocolVersion())) + .put(DRIVER_NAME, driverName().orElse(UNDEFINED)) + .put(DRIVER_VERSION, driverVersion().orElse(UNDEFINED)) + .put(REQUESTS, String.valueOf(requestCount())) + .put(KEYSPACE, keyspace().orElse("")) + .put(SSL, Boolean.toString(isEncrypted())) + .put(CIPHER, sslCipherSuite().orElse(UNDEFINED)) + .put(PROTOCOL, sslProtocol().orElse(UNDEFINED)) + .build(); + } +} http://git-wip-us.apache.org/repos/asf/cassandra/blob/53c71949/src/java/org/apache/cassandra/transport/ConnectionStage.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/transport/ConnectionStage.java b/src/java/org/apache/cassandra/transport/ConnectionStage.java new file mode 100644 index 0000000..128411d --- /dev/null +++ b/src/java/org/apache/cassandra/transport/ConnectionStage.java @@ -0,0 +1,23 @@ +/* + * 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.transport; + +public enum ConnectionStage +{ + ESTABLISHED, AUTHENTICATING, READY +} http://git-wip-us.apache.org/repos/asf/cassandra/blob/53c71949/src/java/org/apache/cassandra/transport/ProtocolVersionTracker.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/transport/ProtocolVersionTracker.java b/src/java/org/apache/cassandra/transport/ProtocolVersionTracker.java index 848917b..72bb901 100644 --- a/src/java/org/apache/cassandra/transport/ProtocolVersionTracker.java +++ b/src/java/org/apache/cassandra/transport/ProtocolVersionTracker.java @@ -15,19 +15,14 @@ * See the License for the specific language governing permissions and * limitations under the License. */ - package org.apache.cassandra.transport; import java.net.InetAddress; +import java.util.ArrayList; import java.util.EnumMap; -import java.util.LinkedHashMap; -import java.util.Map; -import java.util.stream.Collectors; - -import com.google.common.annotations.VisibleForTesting; -import com.google.common.base.Preconditions; -import com.google.common.collect.ImmutableSet; +import java.util.List; +import com.github.benmanes.caffeine.cache.Cache; import com.github.benmanes.caffeine.cache.Caffeine; import com.github.benmanes.caffeine.cache.LoadingCache; import org.apache.cassandra.utils.Clock; @@ -37,17 +32,16 @@ import org.apache.cassandra.utils.Clock; */ public class ProtocolVersionTracker { - public static final int DEFAULT_MAX_CAPACITY = 100; + private static final int DEFAULT_MAX_CAPACITY = 100; - @VisibleForTesting - final EnumMap<ProtocolVersion, LoadingCache<InetAddress, Long>> clientsByProtocolVersion; + private final EnumMap<ProtocolVersion, LoadingCache<InetAddress, Long>> clientsByProtocolVersion; - public ProtocolVersionTracker() + ProtocolVersionTracker() { this(DEFAULT_MAX_CAPACITY); } - public ProtocolVersionTracker(final int capacity) + private ProtocolVersionTracker(int capacity) { clientsByProtocolVersion = new EnumMap<>(ProtocolVersion.class); @@ -58,54 +52,33 @@ public class ProtocolVersionTracker } } - void addConnection(final InetAddress addr, final ProtocolVersion version) + void addConnection(InetAddress addr, ProtocolVersion version) { - if (addr == null || version == null) return; - - LoadingCache<InetAddress, Long> clients = clientsByProtocolVersion.get(version); - clients.put(addr, Clock.instance.currentTimeMillis()); + clientsByProtocolVersion.get(version).put(addr, Clock.instance.currentTimeMillis()); } - public LinkedHashMap<ProtocolVersion, ImmutableSet<ClientIPAndTime>> getAll() + List<ClientStat> getAll() { - LinkedHashMap<ProtocolVersion, ImmutableSet<ClientIPAndTime>> result = new LinkedHashMap<>(); - for (ProtocolVersion version : ProtocolVersion.values()) - { - ImmutableSet.Builder<ClientIPAndTime> ips = ImmutableSet.builder(); - for (Map.Entry<InetAddress, Long> e : clientsByProtocolVersion.get(version).asMap().entrySet()) - ips.add(new ClientIPAndTime(e.getKey(), e.getValue())); - result.put(version, ips.build()); - } + List<ClientStat> result = new ArrayList<>(); + + clientsByProtocolVersion.forEach((version, cache) -> + cache.asMap().forEach((address, lastSeenTime) -> result.add(new ClientStat(address, version, lastSeenTime)))); + return result; } - public void clear() + List<ClientStat> getAll(ProtocolVersion version) { - for (Map.Entry<ProtocolVersion, LoadingCache<InetAddress, Long>> entry : clientsByProtocolVersion.entrySet()) - { - entry.getValue().invalidateAll(); - } - } + List<ClientStat> result = new ArrayList<>(); - public static class ClientIPAndTime - { - final InetAddress inetAddress; - final long lastSeen; + clientsByProtocolVersion.get(version).asMap().forEach((address, lastSeenTime) -> + result.add(new ClientStat(address, version, lastSeenTime))); - public ClientIPAndTime(final InetAddress inetAddress, final long lastSeen) - { - Preconditions.checkNotNull(inetAddress); - this.inetAddress = inetAddress; - this.lastSeen = lastSeen; - } + return result; + } - @Override - public String toString() - { - return "ClientIPAndTime{" + - "inetAddress=" + inetAddress + - ", lastSeen=" + lastSeen + - '}'; - } + public void clear() + { + clientsByProtocolVersion.values().forEach(Cache::invalidateAll); } } http://git-wip-us.apache.org/repos/asf/cassandra/blob/53c71949/src/java/org/apache/cassandra/transport/Server.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/transport/Server.java b/src/java/org/apache/cassandra/transport/Server.java index 996e5bb..8ef137c 100644 --- a/src/java/org/apache/cassandra/transport/Server.java +++ b/src/java/org/apache/cassandra/transport/Server.java @@ -28,9 +28,6 @@ import java.util.concurrent.atomic.AtomicBoolean; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import com.google.common.collect.ImmutableMap; -import com.google.common.collect.ImmutableSet; - import io.netty.bootstrap.ServerBootstrap; import io.netty.buffer.ByteBuf; import io.netty.buffer.ByteBufAllocator; @@ -49,6 +46,7 @@ import io.netty.util.concurrent.EventExecutor; import io.netty.util.concurrent.GlobalEventExecutor; import io.netty.util.internal.logging.InternalLoggerFactory; import io.netty.util.internal.logging.Slf4JLoggerFactory; +import org.apache.cassandra.auth.AuthenticatedUser; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.config.EncryptionOptions; import org.apache.cassandra.db.marshal.AbstractType; @@ -169,63 +167,31 @@ public class Server implements CassandraDaemon.Server isRunning.set(true); } - public int getConnectedClients() + public int countConnectedClients() { - return connectionTracker.getConnectedClients(); + return connectionTracker.countConnectedClients(); } - public Map<String, Integer> getConnectedClientsByUser() + public Map<String, Integer> countConnectedClientsByUser() { - return connectionTracker.getConnectedClientsByUser(); + return connectionTracker.countConnectedClientsByUser(); } - public List<Map<String, String>> getConnectionStates() + public List<ConnectedClient> getConnectedClients() { - List<Map<String, String>> result = new ArrayList<>(); - for(Channel c : connectionTracker.allChannels) + List<ConnectedClient> result = new ArrayList<>(); + for (Channel c : connectionTracker.allChannels) { - Connection connection = c.attr(Connection.attributeKey).get(); - if (connection instanceof ServerConnection) - { - ServerConnection conn = (ServerConnection) connection; - SslHandler sslHandler = conn.channel().pipeline().get(SslHandler.class); - - result.add(new ImmutableMap.Builder<String, String>() - .put("user", conn.getClientState().getUser().getName()) - .put("keyspace", conn.getClientState().getRawKeyspace() == null ? "" : conn.getClientState().getRawKeyspace()) - .put("address", conn.getClientState().getRemoteAddress().toString()) - .put("version", String.valueOf(conn.getVersion().asInt())) - .put("requests", String.valueOf(conn.requests.getCount())) - .put("ssl", Boolean.toString(sslHandler == null)) - .put("cipher", sslHandler != null ? sslHandler.engine().getSession().getCipherSuite() : "undefined") - .put("protocol", sslHandler != null ? sslHandler.engine().getSession().getProtocol() : "undefined") - .put("driverName", conn.getClientState().getDriverName().orElse("undefined")) - .put("driverVersion", conn.getClientState().getDriverVersion().orElse("undefined")) - .build()); - } + Connection conn = c.attr(Connection.attributeKey).get(); + if (conn instanceof ServerConnection) + result.add(new ConnectedClient((ServerConnection) conn)); } return result; } - public List<Map<String, String>> getClientsByProtocolVersion() + public List<ClientStat> recentClientStats() { - LinkedHashMap<ProtocolVersion, ImmutableSet<ProtocolVersionTracker.ClientIPAndTime>> all = connectionTracker.protocolVersionTracker.getAll(); - List<Map<String, String>> result = new ArrayList<>(); - - for (Map.Entry<ProtocolVersion, ImmutableSet<ProtocolVersionTracker.ClientIPAndTime>> entry : all.entrySet()) - { - ProtocolVersion protoVersion = entry.getKey(); - - for (ProtocolVersionTracker.ClientIPAndTime client : entry.getValue()) - { - result.add(new ImmutableMap.Builder<String, String>() - .put("protocolVersion", protoVersion.toString()) - .put("inetAddress", client.inetAddress.toString()) - .put("lastSeenTime", String.valueOf(client.lastSeen)) - .build()); - } - } - return result; + return connectionTracker.protocolVersionTracker.getAll(); } @Override @@ -336,12 +302,12 @@ public class Server implements CassandraDaemon.Server groups.get(event.type).writeAndFlush(new EventMessage(event)); } - public void closeAll() + void closeAll() { allChannels.close().awaitUninterruptibly(); } - public int getConnectedClients() + int countConnectedClients() { /* - When server is running: allChannels contains all clients' connections (channels) @@ -351,16 +317,17 @@ public class Server implements CassandraDaemon.Server return allChannels.size() != 0 ? allChannels.size() - 1 : 0; } - public Map<String, Integer> getConnectedClientsByUser() + Map<String, Integer> countConnectedClientsByUser() { Map<String, Integer> result = new HashMap<>(); - for(Channel c : allChannels) + for (Channel c : allChannels) { Connection connection = c.attr(Connection.attributeKey).get(); if (connection instanceof ServerConnection) { ServerConnection conn = (ServerConnection) connection; - String name = conn.getClientState().getUser().getName(); + AuthenticatedUser user = conn.getClientState().getUser(); + String name = (null != user) ? user.getName() : null; result.put(name, result.getOrDefault(name, 0) + 1); } } http://git-wip-us.apache.org/repos/asf/cassandra/blob/53c71949/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 1ebf81c..d78b7c0 100644 --- a/src/java/org/apache/cassandra/transport/ServerConnection.java +++ b/src/java/org/apache/cassandra/transport/ServerConnection.java @@ -21,20 +21,18 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import io.netty.channel.Channel; +import com.codahale.metrics.Counter; import org.apache.cassandra.auth.IAuthenticator; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.service.ClientState; import org.apache.cassandra.service.QueryState; -import com.codahale.metrics.Counter; - public class ServerConnection extends Connection { - private enum State { UNINITIALIZED, AUTHENTICATION, READY } private volatile IAuthenticator.SaslNegotiator saslNegotiator; private final ClientState clientState; - private volatile State state; + private volatile ConnectionStage stage; public final Counter requests = new Counter(); private final ConcurrentMap<Integer, QueryState> queryStates = new ConcurrentHashMap<>(); @@ -43,7 +41,7 @@ public class ServerConnection extends Connection { super(channel, version, tracker); this.clientState = ClientState.forExternalCalls(channel.remoteAddress()); - this.state = State.UNINITIALIZED; + this.stage = ConnectionStage.ESTABLISHED; } private QueryState getQueryState(int streamId) @@ -64,15 +62,20 @@ public class ServerConnection extends Connection return clientState; } + ConnectionStage stage() + { + return stage; + } + public QueryState validateNewMessage(Message.Type type, ProtocolVersion version, int streamId) { - switch (state) + switch (stage) { - case UNINITIALIZED: + case ESTABLISHED: if (type != Message.Type.STARTUP && type != Message.Type.OPTIONS) throw new ProtocolException(String.format("Unexpected message %s, expecting STARTUP or OPTIONS", type)); break; - case AUTHENTICATION: + case AUTHENTICATING: // Support both SASL auth from protocol v2 and the older style Credentials auth from v1 if (type != Message.Type.AUTH_RESPONSE && type != Message.Type.CREDENTIALS) throw new ProtocolException(String.format("Unexpected message %s, expecting %s", type, version == ProtocolVersion.V1 ? "CREDENTIALS" : "SASL_RESPONSE")); @@ -89,24 +92,24 @@ public class ServerConnection extends Connection public void applyStateTransition(Message.Type requestType, Message.Type responseType) { - switch (state) + switch (stage) { - case UNINITIALIZED: + case ESTABLISHED: if (requestType == Message.Type.STARTUP) { if (responseType == Message.Type.AUTHENTICATE) - state = State.AUTHENTICATION; + stage = ConnectionStage.AUTHENTICATING; else if (responseType == Message.Type.READY) - state = State.READY; + stage = ConnectionStage.READY; } break; - case AUTHENTICATION: + case AUTHENTICATING: // Support both SASL auth from protocol v2 and the older style Credentials auth from v1 assert requestType == Message.Type.AUTH_RESPONSE || requestType == Message.Type.CREDENTIALS; if (responseType == Message.Type.READY || responseType == Message.Type.AUTH_SUCCESS) { - state = State.READY; + stage = ConnectionStage.READY; // we won't use the authenticator again, null it so that it can be GC'd saslNegotiator = null; } http://git-wip-us.apache.org/repos/asf/cassandra/blob/53c71949/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 88c3085..2a20898 100644 --- a/src/java/org/apache/cassandra/transport/messages/AuthResponse.java +++ b/src/java/org/apache/cassandra/transport/messages/AuthResponse.java @@ -25,7 +25,7 @@ import org.apache.cassandra.audit.AuditLogEntryType; import org.apache.cassandra.auth.AuthenticatedUser; import org.apache.cassandra.auth.IAuthenticator; import org.apache.cassandra.exceptions.AuthenticationException; -import org.apache.cassandra.metrics.AuthMetrics; +import org.apache.cassandra.metrics.ClientMetrics; import org.apache.cassandra.service.QueryState; import org.apache.cassandra.transport.*; @@ -80,7 +80,7 @@ public class AuthResponse extends Message.Request { AuthenticatedUser user = negotiator.getAuthenticatedUser(); queryState.getClientState().login(user); - AuthMetrics.instance.markSuccess(); + ClientMetrics.instance.markAuthSuccess(); if (auditLogEnabled) { AuditLogEntry auditEntry = new AuditLogEntry.Builder(queryState.getClientState()) @@ -99,7 +99,7 @@ public class AuthResponse extends Message.Request } catch (AuthenticationException e) { - AuthMetrics.instance.markFailure(); + ClientMetrics.instance.markAuthFailure(); if (auditLogEnabled) { AuditLogEntry auditEntry = new AuditLogEntry.Builder(queryState.getClientState()) http://git-wip-us.apache.org/repos/asf/cassandra/blob/53c71949/test/unit/org/apache/cassandra/transport/ProtocolVersionTrackerTest.java ---------------------------------------------------------------------- diff --git a/test/unit/org/apache/cassandra/transport/ProtocolVersionTrackerTest.java b/test/unit/org/apache/cassandra/transport/ProtocolVersionTrackerTest.java index 6808c0a..91d75b8 100644 --- a/test/unit/org/apache/cassandra/transport/ProtocolVersionTrackerTest.java +++ b/test/unit/org/apache/cassandra/transport/ProtocolVersionTrackerTest.java @@ -20,14 +20,13 @@ package org.apache.cassandra.transport; import java.net.InetAddress; import java.net.UnknownHostException; +import java.util.Collection; import java.util.List; import java.util.stream.Collectors; import java.util.stream.IntStream; -import com.google.common.collect.ImmutableSet; import org.junit.Test; -import static org.apache.cassandra.transport.ProtocolVersionTracker.ClientIPAndTime; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; @@ -45,17 +44,17 @@ public class ProtocolVersionTrackerTest pvt.addConnection(addr, ProtocolVersion.V4); } - ImmutableSet<ClientIPAndTime> clientIPAndTimes1 = pvt.getAll().get(ProtocolVersion.V4); + Collection<ClientStat> clientIPAndTimes1 = pvt.getAll(ProtocolVersion.V4); assertEquals(10, clientIPAndTimes1.size()); Thread.sleep(10); pvt.addConnection(client, ProtocolVersion.V4); - ImmutableSet<ClientIPAndTime> clientIPAndTimes2 = pvt.getAll().get(ProtocolVersion.V4); + Collection<ClientStat> clientIPAndTimes2 = pvt.getAll(ProtocolVersion.V4); assertEquals(10, clientIPAndTimes2.size()); - long ls1 = clientIPAndTimes1.stream().filter(c -> c.inetAddress.equals(client)).findFirst().get().lastSeen; - long ls2 = clientIPAndTimes2.stream().filter(c -> c.inetAddress.equals(client)).findFirst().get().lastSeen; + long ls1 = clientIPAndTimes1.stream().filter(c -> c.remoteAddress.equals(client)).findFirst().get().lastSeenTime; + long ls2 = clientIPAndTimes2.stream().filter(c -> c.remoteAddress.equals(client)).findFirst().get().lastSeenTime; assertTrue(ls2 > ls1); } @@ -75,10 +74,10 @@ public class ProtocolVersionTrackerTest pvt.addConnection(addr, ProtocolVersion.V3); } - assertEquals(5, pvt.getAll().size()); - assertEquals(0, pvt.getAll().get(ProtocolVersion.V2).size()); - assertEquals(7, pvt.getAll().get(ProtocolVersion.V3).size()); - assertEquals(10, pvt.getAll().get(ProtocolVersion.V4).size()); + assertEquals(17, pvt.getAll().size()); + assertEquals(0, pvt.getAll(ProtocolVersion.V2).size()); + assertEquals(7, pvt.getAll(ProtocolVersion.V3).size()); + assertEquals(10, pvt.getAll(ProtocolVersion.V4).size()); } @Test @@ -91,10 +90,10 @@ public class ProtocolVersionTrackerTest pvt.addConnection(addr, ProtocolVersion.V3); } - assertEquals(7, pvt.getAll().get(ProtocolVersion.V3).size()); + assertEquals(7, pvt.getAll(ProtocolVersion.V3).size()); pvt.clear(); - assertEquals(0, pvt.getAll().get(ProtocolVersion.V3).size()); + assertEquals(0, pvt.getAll(ProtocolVersion.V3).size()); } /* Helper */ --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
