diff --git a/src/java/org/apache/cassandra/gms/ApplicationState.java b/src/java/org/apache/cassandra/gms/ApplicationState.java index cbbcf48712..9209c3c7b2 100644 --- a/src/java/org/apache/cassandra/gms/ApplicationState.java +++ b/src/java/org/apache/cassandra/gms/ApplicationState.java @@ -31,6 +31,7 @@ public enum ApplicationState RELEASE_VERSION, REMOVAL_COORDINATOR, INTERNAL_IP, + RPC_ADDRESS, // pad to allow adding new states to existing cluster X1, X2, diff --git a/src/java/org/apache/cassandra/gms/VersionedValue.java b/src/java/org/apache/cassandra/gms/VersionedValue.java index f8a225bf78..93c1faa5e0 100644 --- a/src/java/org/apache/cassandra/gms/VersionedValue.java +++ b/src/java/org/apache/cassandra/gms/VersionedValue.java @@ -21,6 +21,7 @@ package org.apache.cassandra.gms; import java.io.DataInputStream; import java.io.DataOutputStream; import java.io.IOException; +import java.net.InetAddress; import java.util.UUID; import org.apache.cassandra.dht.IPartitioner; @@ -157,6 +158,11 @@ public class VersionedValue implements Comparable return new VersionedValue(rackId); } + public VersionedValue rpcaddress(InetAddress endpoint) + { + return new VersionedValue(endpoint.toString()); + } + public VersionedValue releaseVersion() { return new VersionedValue(FBUtilities.getReleaseVersionString()); diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index 86dd6100fd..ef7e71b33d 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -450,6 +450,12 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe Gossiper.instance.register(this); Gossiper.instance.register(migrationManager); Gossiper.instance.start(SystemTable.incrementAndGetGeneration()); // needed for node-ring gathering. + // add rpc listening info + if (DatabaseDescriptor.getRpcAddress() == null) + Gossiper.instance.addLocalApplicationState(ApplicationState.RPC_ADDRESS, valueFactory.rpcaddress(FBUtilities.getLocalAddress())); + else + Gossiper.instance.addLocalApplicationState(ApplicationState.RPC_ADDRESS, valueFactory.rpcaddress(DatabaseDescriptor.getRpcAddress())); + MessagingService.instance().listen(FBUtilities.getLocalAddress()); StorageLoadBalancer.instance.startBroadcasting(); MigrationManager.passiveAnnounce(DatabaseDescriptor.getDefsVersion()); @@ -574,11 +580,11 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe } /** - * for a keyspace, return the ranges and corresponding hosts for a given keyspace. + * for a keyspace, return the ranges and corresponding RPC addresses for a given keyspace. * @param keyspace * @return */ - public Map> getRangeToEndpointMap(String keyspace) + public Map> getRangeToRpcaddressMap(String keyspace) { // some people just want to get a visual representation of things. Allow null and set it to the first // non-system table. @@ -589,7 +595,15 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe Map> map = new HashMap>(); for (Map.Entry> entry : getRangeToAddressMap(keyspace).entrySet()) { - map.put(entry.getKey(), stringify(entry.getValue())); + List rpcaddrs = new ArrayList(); + for (InetAddress endpoint: entry.getValue()) + { + if (endpoint.equals(FBUtilities.getLocalAddress())) + rpcaddrs.add(DatabaseDescriptor.getRpcAddress().toString()); + else + rpcaddrs.add(Gossiper.instance.getEndpointStateForEndpoint(endpoint).getApplicationState(ApplicationState.RPC_ADDRESS).value); + } + map.put(entry.getKey(), rpcaddrs); } return map; } diff --git a/src/java/org/apache/cassandra/service/StorageServiceMBean.java b/src/java/org/apache/cassandra/service/StorageServiceMBean.java index c78a29672b..1291e69af9 100644 --- a/src/java/org/apache/cassandra/service/StorageServiceMBean.java +++ b/src/java/org/apache/cassandra/service/StorageServiceMBean.java @@ -111,12 +111,12 @@ public interface StorageServiceMBean public String getSavedCachesLocation(); /** - * Retrieve a map of range to end points that describe the ring topology + * Retrieve a map of range to rpc addresses that describe the ring topology * of a Cassandra cluster. * - * @return mapping of ranges to end points + * @return mapping of ranges to rpc addresses */ - public Map> getRangeToEndpointMap(String keyspace); + public Map> getRangeToRpcaddressMap(String keyspace); /** * Retrieve a map of pending ranges to endpoints that describe the ring topology diff --git a/src/java/org/apache/cassandra/thrift/CassandraServer.java b/src/java/org/apache/cassandra/thrift/CassandraServer.java index e50c83a408..7fb9cac324 100644 --- a/src/java/org/apache/cassandra/thrift/CassandraServer.java +++ b/src/java/org/apache/cassandra/thrift/CassandraServer.java @@ -786,7 +786,7 @@ public class CassandraServer implements Cassandra.Iface throw new InvalidRequestException("There is no ring for the keyspace: " + keyspace); List ranges = new ArrayList(); Token.TokenFactory tf = StorageService.getPartitioner().getTokenFactory(); - for (Map.Entry> entry : StorageService.instance.getRangeToEndpointMap(keyspace).entrySet()) + for (Map.Entry> entry : StorageService.instance.getRangeToRpcaddressMap(keyspace).entrySet()) { Range range = entry.getKey(); List endpoints = entry.getValue();