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]

Reply via email to