From 91acae3dfb4a3e8550ca024b91957f695337c6d0 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Tue, 15 Sep 2009 15:47:36 +0000 Subject: [PATCH] Fix deserialization bug in TokenUpdateVerbHandler; change TokenUpdater to update one node at a time, and add bin/tokenupdater as a convenience. Patch by Sammy Yu; reviewed for CASSANDRA-363 by jbellis git-svn-id: https://svn.apache.org/repos/asf/incubator/cassandra/trunk@815373 13f79535-47bb-0310-9956-ffa450edef68 --- bin/tokenupdater | 49 ++++++++++ .../cassandra/net/MessagingService.java | 2 + .../service/TokenUpdateVerbHandler.java | 22 +++-- .../tools/TokenUpdateVerbHandler.java | 94 ------------------- .../apache/cassandra/tools/TokenUpdater.java | 51 +++++----- 5 files changed, 93 insertions(+), 125 deletions(-) create mode 100644 bin/tokenupdater delete mode 100644 src/java/org/apache/cassandra/tools/TokenUpdateVerbHandler.java diff --git a/bin/tokenupdater b/bin/tokenupdater new file mode 100644 index 0000000000..12118e9791 --- /dev/null +++ b/bin/tokenupdater @@ -0,0 +1,49 @@ +#!/bin/sh +# 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. + + +if [ "x$CASSANDRA_INCLUDE" = "x" ]; then + for include in /usr/share/cassandra/cassandra.in.sh \ + /usr/local/share/cassandra/cassandra.in.sh \ + /opt/cassandra/cassandra.in.sh \ + `dirname $0`/cassandra.in.sh; do + if [ -r $include ]; then + . $include + break + fi + done +elif [ -r $CASSANDRA_INCLUDE ]; then + . $CASSANDRA_INCLUDE +fi + +if [ -z $CASSANDRA_CONF -o -z $CLASSPATH ]; then + echo "You must set the CASSANDRA_CONF and CLASSPATH vars" >&2 + exit 1 +fi + +# Special-case path variables. +case "`uname`" in + CYGWIN*) + CLASSPATH=`cygpath -p -w "$CLASSPATH"` + CASSANDRA_CONF=`cygpath -p -w "$CASSANDRA_CONF"` + ;; +esac + +java -cp $CLASSPATH -Dstorage-config=$CASSANDRA_CONF \ + org.apache.cassandra.tools.TokenUpdater $@ + +# vi:ai sw=4 ts=4 tw=0 et diff --git a/src/java/org/apache/cassandra/net/MessagingService.java b/src/java/org/apache/cassandra/net/MessagingService.java index 21f97d2189..d38b055c59 100644 --- a/src/java/org/apache/cassandra/net/MessagingService.java +++ b/src/java/org/apache/cassandra/net/MessagingService.java @@ -514,6 +514,8 @@ public class MessagingService implements IMessagingService messageDeserializerExecutor_.shutdownNow(); streamExecutor_.shutdownNow(); + StageManager.shutdown(); + /* shut down the cachetables */ taskCompletionMap_.shutdown(); callbackMap_.shutdown(); diff --git a/src/java/org/apache/cassandra/service/TokenUpdateVerbHandler.java b/src/java/org/apache/cassandra/service/TokenUpdateVerbHandler.java index be5410cfaf..8fa152baaa 100644 --- a/src/java/org/apache/cassandra/service/TokenUpdateVerbHandler.java +++ b/src/java/org/apache/cassandra/service/TokenUpdateVerbHandler.java @@ -23,9 +23,9 @@ import java.io.IOException; import org.apache.log4j.Logger; import org.apache.cassandra.dht.Token; +import org.apache.cassandra.io.DataInputBuffer; import org.apache.cassandra.net.IVerbHandler; import org.apache.cassandra.net.Message; -import org.apache.cassandra.utils.LogUtil; public class TokenUpdateVerbHandler implements IVerbHandler { @@ -33,18 +33,20 @@ public class TokenUpdateVerbHandler implements IVerbHandler public void doVerb(Message message) { - byte[] body = message.getMessageBody(); - Token token = StorageService.getPartitioner().getTokenFactory().fromByteArray(body); + byte[] body = message.getMessageBody(); + DataInputBuffer bufIn = new DataInputBuffer(); + bufIn.reset(body, body.length); try { - logger_.info("Updating the token to [" + token + "]"); - StorageService.instance().updateToken(token); + /* Deserialize to get the token for this endpoint. */ + Token token = Token.serializer().deserialize(bufIn); + logger_.info("Updating the token to [" + token + "]"); + StorageService.instance().updateToken(token); + } + catch (IOException ex) + { + throw new RuntimeException(ex); } - catch( IOException ex ) - { - if (logger_.isDebugEnabled()) - logger_.debug(LogUtil.throwableToString(ex)); - } } } diff --git a/src/java/org/apache/cassandra/tools/TokenUpdateVerbHandler.java b/src/java/org/apache/cassandra/tools/TokenUpdateVerbHandler.java deleted file mode 100644 index 236fc83b38..0000000000 --- a/src/java/org/apache/cassandra/tools/TokenUpdateVerbHandler.java +++ /dev/null @@ -1,94 +0,0 @@ -/** - * 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.tools; - -import java.io.ByteArrayOutputStream; -import java.io.DataOutputStream; -import java.io.IOException; -import java.util.Map; -import java.util.Set; - -import org.apache.log4j.Logger; - -import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.dht.IPartitioner; -import org.apache.cassandra.dht.Token; -import org.apache.cassandra.io.DataInputBuffer; -import org.apache.cassandra.net.EndPoint; -import org.apache.cassandra.net.IVerbHandler; -import org.apache.cassandra.net.Message; -import org.apache.cassandra.net.MessagingService; -import org.apache.cassandra.service.StorageService; -import org.apache.cassandra.utils.LogUtil; - -public class TokenUpdateVerbHandler implements IVerbHandler -{ - private static Logger logger_ = Logger.getLogger(TokenUpdateVerbHandler.class); - - public void doVerb(Message message) - { - byte[] body = message.getMessageBody(); - - try - { - DataInputBuffer bufIn = new DataInputBuffer(); - bufIn.reset(body, body.length); - /* Deserialize to get the token for this endpoint. */ - Token token = Token.serializer().deserialize(bufIn); - - logger_.info("Updating the token to [" + token + "]"); - StorageService.instance().updateToken(token); - - /* Get the headers for this message */ - Map headers = message.getHeaders(); - headers.remove( StorageService.getLocalStorageEndPoint().getHost() ); - if (logger_.isDebugEnabled()) - logger_.debug("Number of nodes in the header " + headers.size()); - Set nodes = headers.keySet(); - - IPartitioner p = StorageService.getPartitioner(); - for ( String node : nodes ) - { - if (logger_.isDebugEnabled()) - logger_.debug("Processing node " + node); - byte[] bytes = headers.remove(node); - /* Send a message to this node to update its token to the one retrieved. */ - EndPoint target = new EndPoint(node, DatabaseDescriptor.getStoragePort()); - token = p.getTokenFactory().fromByteArray(bytes); - - /* Reset the new Message */ - ByteArrayOutputStream bos = new ByteArrayOutputStream(); - DataOutputStream dos = new DataOutputStream(bos); - Token.serializer().serialize(token, dos); - message.setMessageBody(bos.toByteArray()); - - if (logger_.isDebugEnabled()) - logger_.debug("Sending a token update message to " + target + " to update it to " + token); - MessagingService.getMessagingInstance().sendOneWay(message, target); - break; - } - } - catch( IOException ex ) - { - if (logger_.isDebugEnabled()) - logger_.debug(LogUtil.throwableToString(ex)); - } - } - -} diff --git a/src/java/org/apache/cassandra/tools/TokenUpdater.java b/src/java/org/apache/cassandra/tools/TokenUpdater.java index 378e2645aa..56abe51084 100644 --- a/src/java/org/apache/cassandra/tools/TokenUpdater.java +++ b/src/java/org/apache/cassandra/tools/TokenUpdater.java @@ -23,58 +23,67 @@ import java.io.ByteArrayOutputStream; import java.io.DataOutputStream; import java.io.FileInputStream; import java.io.InputStreamReader; +import java.io.ByteArrayInputStream; +import java.io.DataInputStream; import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.dht.Token; import org.apache.cassandra.net.EndPoint; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessagingService; +import org.apache.cassandra.net.SelectorManager; import org.apache.cassandra.service.StorageService; import org.apache.cassandra.utils.FBUtilities; +import org.apache.cassandra.utils.FileUtils; public class TokenUpdater { private static final int port_ = 7000; private static final long waitTime_ = 10000; - + public static void main(String[] args) throws Throwable { - if ( args.length != 3 ) + if (args.length < 2) { - System.out.println("Usage : java org.apache.cassandra.tools.TokenUpdater "); + System.out.println("Usage : java org.apache.cassandra.tools.TokenUpdater "); System.exit(1); } - + + Thread selectorThread = SelectorManager.getSelectorManager(); + selectorThread.setDaemon(true); + selectorThread.start(); + String ipPort = args[0]; IPartitioner p = StorageService.getPartitioner(); Token token = p.getTokenFactory().fromString(args[1]); - String file = args[2]; - + System.out.println("Partitioner is " + p.getClass() + ", token is: " + token); + System.out.println(p.getTokenFactory().getClass()); + String[] ipPortPair = ipPort.split(":"); - EndPoint target = new EndPoint(ipPortPair[0], Integer.valueOf(ipPortPair[1])); + int port = 7000; + if (ipPortPair.length > 1) + { + port = Integer.valueOf(ipPortPair[1]); + } + + EndPoint target = new EndPoint(ipPortPair[0], port); ByteArrayOutputStream bos = new ByteArrayOutputStream(); DataOutputStream dos = new DataOutputStream(bos); Token.serializer().serialize(token, dos); /* Construct the token update message to be sent */ - Message tokenUpdateMessage = new Message( new EndPoint(FBUtilities.getHostAddress(), port_), "", StorageService.tokenVerbHandler_, bos.toByteArray() ); - - BufferedReader bufReader = new BufferedReader( new InputStreamReader( new FileInputStream(file) ) ); - String line = null; - - while ( ( line = bufReader.readLine() ) != null ) - { - String[] nodeTokenPair = line.split(" "); - /* Add the node and the token pair into the header of this message. */ - Token nodeToken = p.getTokenFactory().fromString(nodeTokenPair[1]); - tokenUpdateMessage.addHeader(nodeTokenPair[0], p.getTokenFactory().toByteArray(nodeToken)); - } - + Message tokenUpdateMessage = new Message(new EndPoint(FBUtilities.getHostAddress(), port_), + "", + StorageService.tokenVerbHandler_, + bos.toByteArray()); + System.out.println("Sending a token update message to " + target); MessagingService.getMessagingInstance().sendOneWay(tokenUpdateMessage, target); Thread.sleep(TokenUpdater.waitTime_); System.out.println("Done sending the update message"); - } + MessagingService.shutdown(); + FileUtils.shutdown(); + } }