Updates for allowing commands to be cancelled, within single processes is 
complete.  Now need to add a spray to the cluster when cancel is called as well 
as forcing indexreaders to exit.  Exception handling may need work.


Project: http://git-wip-us.apache.org/repos/asf/incubator-blur/repo
Commit: http://git-wip-us.apache.org/repos/asf/incubator-blur/commit/1bbb14b2
Tree: http://git-wip-us.apache.org/repos/asf/incubator-blur/tree/1bbb14b2
Diff: http://git-wip-us.apache.org/repos/asf/incubator-blur/diff/1bbb14b2

Branch: refs/heads/master
Commit: 1bbb14b2d76d233a99774a8a20df702b7b7d94de
Parents: 0238b4a
Author: Aaron McCurry <[email protected]>
Authored: Fri Oct 31 10:37:07 2014 -0400
Committer: Aaron McCurry <[email protected]>
Committed: Fri Oct 31 10:37:07 2014 -0400

----------------------------------------------------------------------
 .../apache/blur/command/ArgumentOverlay.java    |  20 ++++
 .../apache/blur/command/BaseCommandManager.java | 109 ++++++++++++-------
 .../java/org/apache/blur/command/Command.java   |  11 ++
 .../org/apache/blur/command/CommandRunner.java  |   4 +-
 .../blur/command/ControllerClusterContext.java  |   8 +-
 .../blur/command/ControllerCommandManager.java  |   2 +-
 .../org/apache/blur/command/ResponseFuture.java |  16 ++-
 .../blur/command/ShardCommandManager.java       |  17 +--
 .../apache/blur/command/TimeoutException.java   |  10 +-
 .../apache/blur/server/FilteredBlurServer.java  |  22 ++--
 .../blur/thrift/BlurControllerServer.java       |  15 ++-
 .../org/apache/blur/thrift/BlurShardServer.java |  18 +--
 .../blur/command/ShardCommandManagerTest.java   |  81 +++++++++++---
 13 files changed, 225 insertions(+), 108 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/ArgumentOverlay.java
----------------------------------------------------------------------
diff --git 
a/blur-core/src/main/java/org/apache/blur/command/ArgumentOverlay.java 
b/blur-core/src/main/java/org/apache/blur/command/ArgumentOverlay.java
index e051b93..493f1f8 100644
--- a/blur-core/src/main/java/org/apache/blur/command/ArgumentOverlay.java
+++ b/blur-core/src/main/java/org/apache/blur/command/ArgumentOverlay.java
@@ -18,12 +18,17 @@ package org.apache.blur.command;
 
 import java.lang.reflect.Field;
 import java.util.Map;
+import java.util.UUID;
 
 import org.apache.blur.command.annotation.OptionalArgument;
 import org.apache.blur.command.annotation.RequiredArgument;
+import org.apache.blur.log.Log;
+import org.apache.blur.log.LogFactory;
 
 public class ArgumentOverlay {
 
+  private static final Log LOG = LogFactory.getLog(ArgumentOverlay.class);
+
   private final Map<String, ? extends Object> _args;
 
   public ArgumentOverlay(BlurObject args, BlurObjectSerDe serDe) {
@@ -33,14 +38,29 @@ public class ArgumentOverlay {
   public <T> Command<T> setup(Command<T> command) {
     Class<?> clazz = command.getClass();
     setupInternal(clazz, command);
+    return setupCommandExecutionId(command);
+  }
+
+  private <T> Command<T> setupCommandExecutionId(Command<T> command) {
+    String commandExecutionId = command.getCommandExecutionId();
+    if (commandExecutionId == null) {
+      commandExecutionId = UUID.randomUUID().toString();
+      LOG.info("Command execution id [{0}] has been assigned to [{1}]", 
commandExecutionId, command);
+      command.setCommandExecutionId(commandExecutionId);
+    }
     return command;
   }
 
   private void setupInternal(Class<?> clazz, Command<?> command) {
     if (clazz.equals(Command.class)) {
+      mapValuesToFields(clazz, command);
       return;
     }
     setupInternal(clazz.getSuperclass(), command);
+    mapValuesToFields(clazz, command);
+  }
+
+  private void mapValuesToFields(Class<?> clazz, Command<?> command) {
     Field[] declaredFields = clazz.getDeclaredFields();
     for (Field field : declaredFields) {
       RequiredArgument requiredArgument = 
field.getAnnotation(RequiredArgument.class);

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/BaseCommandManager.java
----------------------------------------------------------------------
diff --git 
a/blur-core/src/main/java/org/apache/blur/command/BaseCommandManager.java 
b/blur-core/src/main/java/org/apache/blur/command/BaseCommandManager.java
index ba6d4be..6ccded6 100644
--- a/blur-core/src/main/java/org/apache/blur/command/BaseCommandManager.java
+++ b/blur-core/src/main/java/org/apache/blur/command/BaseCommandManager.java
@@ -16,6 +16,7 @@ import java.util.List;
 import java.util.Map;
 import java.util.Map.Entry;
 import java.util.Properties;
+import java.util.Random;
 import java.util.Set;
 import java.util.Timer;
 import java.util.TimerTask;
@@ -66,13 +67,15 @@ public abstract class BaseCommandManager implements 
Closeable {
   private static final String 
META_INF_SERVICES_ORG_APACHE_BLUR_COMMAND_COMMANDS = 
"META-INF/services/org.apache.blur.command.Commands";
   private static final Log LOG = LogFactory.getLog(BaseCommandManager.class);
 
-  private final ExecutorService _executorService;
+  private final ExecutorService _executorServiceWorker;
   private final ExecutorService _executorServiceDriver;
+  private final Random _random = new Random();
 
   protected final Map<String, BigInteger> _commandLoadTime = new 
ConcurrentHashMap<String, BigInteger>();
   protected final Map<String, Command<?>> _command = new 
ConcurrentHashMap<String, Command<?>>();
   protected final Map<Class<? extends Command<?>>, String> _commandNameLookup 
= new ConcurrentHashMap<Class<? extends Command<?>>, String>();
-  protected final ConcurrentMap<ExecutionId, ResponseFuture> _runningMap = new 
MapMaker().makeMap();
+  protected final ConcurrentMap<Long, ResponseFuture<?>> _driverRunningMap = 
new MapMaker().makeMap();
+  protected final ConcurrentMap<Long, ResponseFuture<?>> _workerRunningMap = 
new MapMaker().makeMap();
   protected final long _connectionTimeout;
   protected final File _tmpPath;
   protected final String _commandPath;
@@ -89,11 +92,12 @@ public abstract class BaseCommandManager implements 
Closeable {
     lookForCommandsToRegisterInClassPath();
     _tmpPath = tmpPath;
     _commandPath = commandPath;
-    _executorService = Executors.newThreadPool("command-worker-", 
workerThreadCount);
+    _executorServiceWorker = Executors.newThreadPool("command-worker-", 
workerThreadCount);
     _executorServiceDriver = Executors.newThreadPool("command-driver-", 
driverThreadCount);
     _connectionTimeout = connectionTimeout / 2;
     _timer = new Timer("BaseCommandManager-Timer", true);
-    _timer.schedule(getTimerTaskForRemovalOfOldCommands(), _pollingPeriod, 
_pollingPeriod);
+    _timer.schedule(getTimerTaskForRemovalOfOldCommands(_driverRunningMap), 
_pollingPeriod, _pollingPeriod);
+    _timer.schedule(getTimerTaskForRemovalOfOldCommands(_workerRunningMap), 
_pollingPeriod, _pollingPeriod);
     if (_tmpPath == null || _commandPath == null) {
       LOG.info("Tmp Path [{0}] or Command Path [{1}] is null so the automatic 
command reload will be disabled.",
           _tmpPath, _commandPath);
@@ -103,30 +107,45 @@ public abstract class BaseCommandManager implements 
Closeable {
     }
   }
 
-  public void cancel(ExecutionId executionId) {
-    ResponseFuture responseFuture = _runningMap.get(executionId);
-    responseFuture.cancel(true);
+  public void cancelCommand(String commandExecutionId) {
+    LOG.info("Trying to cancel command [{0}]", commandExecutionId);
+    cancelAllExecuting(commandExecutionId, _workerRunningMap);
+    cancelAllExecuting(commandExecutionId, _driverRunningMap);
   }
 
-  public List<String> commandStatusList(CommandStatusStateEnum commandStatus) {
-    throw new RuntimeException("Not implemented.");
+  private void cancelAllExecuting(String commandExecutionId, 
ConcurrentMap<Long, ResponseFuture<?>> runningMap) {
+    for (Entry<Long, ResponseFuture<?>> e : runningMap.entrySet()) {
+      Long instanceExecutionId = e.getKey();
+      ResponseFuture<?> future = e.getValue();
+      Command<?> commandExecuting = future.getCommandExecuting();
+      if (commandExecuting.getCommandExecutionId().equals(commandExecutionId)) 
{
+        LOG.info("Canceling Command with executing id [{0}] command [{1}]", 
instanceExecutionId, commandExecuting);
+        future.cancel(true);
+      }
+    }
   }
 
-  public void cancel(String executionId) {
-    cancel(new ExecutionId(executionId));
+  public List<String> commandStatusList(CommandStatusStateEnum commandStatus) {
+    throw new RuntimeException("Not implemented.");
   }
 
-  private TimerTask getTimerTaskForRemovalOfOldCommands() {
+  private TimerTask getTimerTaskForRemovalOfOldCommands(final Map<Long, 
ResponseFuture<?>> runningMap) {
     return new TimerTask() {
       @Override
       public void run() {
-        Set<Entry<ExecutionId, ResponseFuture>> entrySet = 
_runningMap.entrySet();
-        for (Entry<ExecutionId, ResponseFuture> e : entrySet) {
-          ExecutionId executionId = e.getKey();
-          ResponseFuture responseFuture = e.getValue();
+        Set<Entry<Long, ResponseFuture<?>>> entrySet = runningMap.entrySet();
+        for (Entry<Long, ResponseFuture<?>> e : entrySet) {
+          Long instanceExecutionId = e.getKey();
+          ResponseFuture<?> responseFuture = e.getValue();
           if (!responseFuture.isRunning() && responseFuture.hasExpired()) {
-            LOG.info("Removing old execution id [{0}]", executionId);
-            _runningMap.remove(executionId);
+            Command<?> commandExecuting = responseFuture.getCommandExecuting();
+            String commandExecutionId = null;
+            if (commandExecuting != null) {
+              commandExecutionId = commandExecuting.getCommandExecutionId();
+            }
+            LOG.info("Removing old execution instance id [{0}] with command 
execution id of [{1}]",
+                instanceExecutionId, commandExecutionId);
+            runningMap.remove(instanceExecutionId);
           }
         }
       }
@@ -284,10 +303,11 @@ public abstract class BaseCommandManager implements 
Closeable {
     }
   }
 
-  public Response reconnect(ExecutionId executionId) throws IOException, 
TimeoutException {
-    Future<Response> future = _runningMap.get(executionId);
+  @SuppressWarnings("unchecked")
+  public Response reconnect(Long instanceExecutionId) throws IOException, 
TimeoutException {
+    Future<Response> future = (Future<Response>) 
_driverRunningMap.get(instanceExecutionId);
     if (future == null) {
-      throw new IOException("Command id [" + executionId + "] did not find any 
executing commands.");
+      throw new IOException("Execution instance id [" + instanceExecutionId + 
"] did not find any executing commands.");
     }
     try {
       return future.get(_connectionTimeout, TimeUnit.MILLISECONDS);
@@ -298,18 +318,17 @@ public abstract class BaseCommandManager implements 
Closeable {
     } catch (ExecutionException e) {
       throw new IOException(e.getCause());
     } catch (java.util.concurrent.TimeoutException e) {
-      LOG.info("Timeout of command [{0}]", executionId);
-      throw new TimeoutException(executionId);
+      LOG.info("Timeout of command [{0}]", instanceExecutionId);
+      throw new TimeoutException(instanceExecutionId);
     }
   }
 
-  protected Response submitDriverCallable(Callable<Response> callable) throws 
IOException, TimeoutException,
-      ExceptionCollector {
-    ExecutionContext executionContext = ExecutionContext.create();
-    Future<Response> future = 
_executorServiceDriver.submit(executionContext.wrapCallable(callable));
-    executionContext.registerDriverFuture(future);
-    ExecutionId executionId = executionContext.getExecutionId();
-    _runningMap.put(executionId, new 
ResponseFuture(_runningCacheTombstoneTime, future));
+  protected Response submitDriverCallable(Callable<Response> callable, 
Command<?> commandExecuting) throws IOException,
+      TimeoutException, ExceptionCollector {
+    Future<Response> future = _executorServiceDriver.submit(callable);
+    Long instanceExecutionId = getInstanceExecutionId();
+    _driverRunningMap.put(instanceExecutionId, new 
ResponseFuture<Response>(_runningCacheTombstoneTime, future,
+        commandExecuting));
     try {
       return future.get(_connectionTimeout, TimeUnit.MILLISECONDS);
     } catch (CancellationException e) {
@@ -323,21 +342,37 @@ public abstract class BaseCommandManager implements 
Closeable {
       }
       throw new IOException(cause);
     } catch (java.util.concurrent.TimeoutException e) {
-      LOG.info("Timeout of command [{0}]", executionId);
-      throw new TimeoutException(executionId);
+      LOG.info("Timeout of command [{0}]", instanceExecutionId);
+      throw new TimeoutException(instanceExecutionId);
+    }
+  }
+
+  private Long getInstanceExecutionId() {
+    synchronized (_random) {
+      while (true) {
+        Long id = _random.nextLong();
+        if (_driverRunningMap.containsKey(id)) {
+          continue;
+        }
+        if (_workerRunningMap.containsKey(id)) {
+          continue;
+        }
+        return id;
+      }
     }
   }
 
-  protected <T> Future<T> submitToExecutorService(Callable<T> callable) {
-    ExecutionContext executionContext = ExecutionContext.get();
-    Future<T> future = 
_executorService.submit(executionContext.wrapCallable(callable));
-    executionContext.registerFuture(future);
+  protected <T> Future<T> submitToExecutorService(Callable<T> callable, 
Command<?> commandExecuting) {
+    Future<T> future = _executorServiceWorker.submit(callable);
+    Long instanceExecutionId = getInstanceExecutionId();
+    _workerRunningMap.put(instanceExecutionId, new 
ResponseFuture<T>(_runningCacheTombstoneTime, future,
+        commandExecuting));
     return future;
   }
 
   @Override
   public void close() throws IOException {
-    _executorService.shutdownNow();
+    _executorServiceWorker.shutdownNow();
     _executorServiceDriver.shutdownNow();
     if (_timer != null) {
       _timer.cancel();

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/Command.java
----------------------------------------------------------------------
diff --git a/blur-core/src/main/java/org/apache/blur/command/Command.java 
b/blur-core/src/main/java/org/apache/blur/command/Command.java
index 89b6e0e..9cf2719 100644
--- a/blur-core/src/main/java/org/apache/blur/command/Command.java
+++ b/blur-core/src/main/java/org/apache/blur/command/Command.java
@@ -30,6 +30,9 @@ import org.apache.blur.thrift.generated.Blur.Iface;
 
 public abstract class Command<R> implements Cloneable {
 
+  @OptionalArgument("The ")
+  private String commandExecutionId;
+
   public abstract String getName();
 
   public abstract R run() throws IOException;
@@ -117,4 +120,12 @@ public abstract class Command<R> implements Cloneable {
     return builder.toString();
   }
 
+  public String getCommandExecutionId() {
+    return commandExecutionId;
+  }
+
+  public void setCommandExecutionId(String commandExecutionId) {
+    this.commandExecutionId = commandExecutionId;
+  }
+
 }

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/CommandRunner.java
----------------------------------------------------------------------
diff --git a/blur-core/src/main/java/org/apache/blur/command/CommandRunner.java 
b/blur-core/src/main/java/org/apache/blur/command/CommandRunner.java
index f88d6c2..38c38aa 100644
--- a/blur-core/src/main/java/org/apache/blur/command/CommandRunner.java
+++ b/blur-core/src/main/java/org/apache/blur/command/CommandRunner.java
@@ -209,7 +209,7 @@ public class CommandRunner {
       ClientPool clientPool = BlurClientManager.getClientPool();
       Client client = clientPool.getClient(connection);
       try {
-        String executionId = null;
+        Long executionId = null;
         Response response;
         INNER: while (true) {
           try {
@@ -220,7 +220,7 @@ public class CommandRunner {
             }
             break INNER;
           } catch (TimeoutException te) {
-            executionId = te.getExecutionId();
+            executionId = te.getInstanceExecutionId();
           }
         }
         return CommandUtil.fromThriftResponseToObject(response);

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/ControllerClusterContext.java
----------------------------------------------------------------------
diff --git 
a/blur-core/src/main/java/org/apache/blur/command/ControllerClusterContext.java 
b/blur-core/src/main/java/org/apache/blur/command/ControllerClusterContext.java
index 29c4c70..13b1589 100644
--- 
a/blur-core/src/main/java/org/apache/blur/command/ControllerClusterContext.java
+++ 
b/blur-core/src/main/java/org/apache/blur/command/ControllerClusterContext.java
@@ -126,7 +126,7 @@ public class ControllerClusterContext extends 
ClusterContext implements Closeabl
           Map<Shard, Object> shardToValue = 
CommandUtil.fromThriftSupportedObjects(shardToThriftValue, _serDe);
           return (Map<Shard, T>) shardToValue;
         }
-      });
+      }, command);
       for (Shard shard : getShardsOnServer(server, tables, shards)) {
         futureMap.put(shard, new ShardResultFuture<T>(shard, future));
       }
@@ -149,7 +149,7 @@ public class ControllerClusterContext extends 
ClusterContext implements Closeabl
   protected static Response waitForResponse(Client client, Command<?> command, 
Arguments arguments) throws TException {
     // TODO This should likely be changed to run of a AtomicBoolean used for
     // the status of commands.
-    String executionId = null;
+    Long executionId = null;
     while (true) {
       try {
         if (executionId == null) {
@@ -160,7 +160,7 @@ public class ControllerClusterContext extends 
ClusterContext implements Closeabl
       } catch (BlurException e) {
         throw e;
       } catch (TimeoutException e) {
-        executionId = e.getExecutionId();
+        executionId = e.getInstanceExecutionId();
         LOG.info("Execution fetch timed out, reconnecting using [{0}].", 
executionId);
       } catch (TException e) {
         throw e;
@@ -212,7 +212,7 @@ public class ControllerClusterContext extends 
ClusterContext implements Closeabl
           Object thriftObject = CommandUtil.toObject(valueObject);
           return (T) _serDe.fromSupportedThriftObject(thriftObject);
         }
-      });
+      }, command);
       futureMap.put(server, future);
     }
     return futureMap;

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/ControllerCommandManager.java
----------------------------------------------------------------------
diff --git 
a/blur-core/src/main/java/org/apache/blur/command/ControllerCommandManager.java 
b/blur-core/src/main/java/org/apache/blur/command/ControllerCommandManager.java
index 31901f8..386bae8 100644
--- 
a/blur-core/src/main/java/org/apache/blur/command/ControllerCommandManager.java
+++ 
b/blur-core/src/main/java/org/apache/blur/command/ControllerCommandManager.java
@@ -77,7 +77,7 @@ public class ControllerCommandManager extends 
BaseCommandManager {
         throw new IOException("Command type of [" + command.getClass() + "] 
not supported.");
       }
 
-    });
+    }, command);
   }
 
   private CombiningContext getCombiningContext(final TableContextFactory 
tableContextFactory) {

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/ResponseFuture.java
----------------------------------------------------------------------
diff --git 
a/blur-core/src/main/java/org/apache/blur/command/ResponseFuture.java 
b/blur-core/src/main/java/org/apache/blur/command/ResponseFuture.java
index ca75b3e..eb9b30d 100644
--- a/blur-core/src/main/java/org/apache/blur/command/ResponseFuture.java
+++ b/blur-core/src/main/java/org/apache/blur/command/ResponseFuture.java
@@ -22,15 +22,21 @@ import java.util.concurrent.TimeUnit;
 import java.util.concurrent.TimeoutException;
 import java.util.concurrent.atomic.AtomicLong;
 
-public class ResponseFuture implements Future<Response> {
+public class ResponseFuture<T> implements Future<T> {
 
-  private final Future<Response> _future;
+  private final Future<T> _future;
   private final AtomicLong _timeWhenNotRunningObserved = new AtomicLong();
   private final long _tombstone;
+  private final Command<?> _commandExecuting;
 
-  public ResponseFuture(long tombstone, Future<Response> future) {
+  public ResponseFuture(long tombstone, Future<T> future, Command<?> 
commandExecuting) {
     _tombstone = tombstone;
     _future = future;
+    _commandExecuting = commandExecuting;
+  }
+
+  public Command<?> getCommandExecuting() {
+    return _commandExecuting;
   }
 
   public boolean cancel(boolean mayInterruptIfRunning) {
@@ -45,11 +51,11 @@ public class ResponseFuture implements Future<Response> {
     return _future.isDone();
   }
 
-  public Response get() throws InterruptedException, ExecutionException {
+  public T get() throws InterruptedException, ExecutionException {
     return _future.get();
   }
 
-  public Response get(long timeout, TimeUnit unit) throws 
InterruptedException, ExecutionException, TimeoutException {
+  public T get(long timeout, TimeUnit unit) throws InterruptedException, 
ExecutionException, TimeoutException {
     return _future.get(timeout, unit);
   }
 

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/ShardCommandManager.java
----------------------------------------------------------------------
diff --git 
a/blur-core/src/main/java/org/apache/blur/command/ShardCommandManager.java 
b/blur-core/src/main/java/org/apache/blur/command/ShardCommandManager.java
index ec3b6ab..cbf218f 100644
--- a/blur-core/src/main/java/org/apache/blur/command/ShardCommandManager.java
+++ b/blur-core/src/main/java/org/apache/blur/command/ShardCommandManager.java
@@ -51,10 +51,10 @@ public class ShardCommandManager extends BaseCommandManager 
{
   public Response execute(final TableContextFactory tableContextFactory, final 
String commandName,
       final ArgumentOverlay argumentOverlay) throws IOException, 
TimeoutException, ExceptionCollector {
     final ShardServerContext shardServerContext = getShardServerContext();
+    final Command<?> command = getCommandObject(commandName, argumentOverlay);
     Callable<Response> callable = new Callable<Response>() {
       @Override
       public Response call() throws Exception {
-        Command<?> command = getCommandObject(commandName, argumentOverlay);
         if (command == null) {
           throw new IOException("Command with name [" + commandName + "] not 
found.");
         }
@@ -80,7 +80,7 @@ public class ShardCommandManager extends BaseCommandManager {
         };
       }
     };
-    return submitDriverCallable(callable);
+    return submitDriverCallable(callable, command);
   }
 
   private ShardServerContext getShardServerContext() {
@@ -137,16 +137,17 @@ public class ShardCommandManager extends 
BaseCommandManager {
         }
         final BlurIndex blurIndex = e.getValue();
         Callable<Object> callable;
-        if (command instanceof IndexRead) {
-          final IndexRead<?> readCommand = (IndexRead<?>) command.clone();
+        Command<?> clone = command.clone();
+        if (clone instanceof IndexRead) {
+          final IndexRead<?> readCommand = (IndexRead<?>) clone;
           callable = getCallable(shardServerContext, tableContextFactory, 
table, shard, blurIndex, readCommand);
-        } else if (command instanceof ServerRead) {
-          final ServerRead<?, ?> readCombiningCommand = (ServerRead<?, ?>) 
command.clone();
+        } else if (clone instanceof ServerRead) {
+          final ServerRead<?, ?> readCombiningCommand = (ServerRead<?, ?>) 
clone;
           callable = getCallable(shardServerContext, tableContextFactory, 
table, shard, blurIndex, readCombiningCommand);
         } else {
-          throw new IOException("Command type of [" + command.getClass() + "] 
not supported.");
+          throw new IOException("Command type of [" + clone.getClass() + "] 
not supported.");
         }
-        Future<Object> future = submitToExecutorService(callable);
+        Future<Object> future = submitToExecutorService(callable, clone);
         futureMap.put(shard, future);
       }
     }

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/TimeoutException.java
----------------------------------------------------------------------
diff --git 
a/blur-core/src/main/java/org/apache/blur/command/TimeoutException.java 
b/blur-core/src/main/java/org/apache/blur/command/TimeoutException.java
index a0b091f..9f85bb2 100644
--- a/blur-core/src/main/java/org/apache/blur/command/TimeoutException.java
+++ b/blur-core/src/main/java/org/apache/blur/command/TimeoutException.java
@@ -19,14 +19,14 @@ package org.apache.blur.command;
 @SuppressWarnings("serial")
 public class TimeoutException extends Exception {
 
-  private final ExecutionId _executionId;
+  private final Long _executionInstanceId;
 
-  public TimeoutException(ExecutionId executionId) {
-    _executionId = executionId;
+  public TimeoutException(Long executionInstanceId) {
+    _executionInstanceId = executionInstanceId;
   }
 
-  public ExecutionId getExecutionId() {
-    return _executionId;
+  public Long getInstanceExecutionId() {
+    return _executionInstanceId;
   }
 
 }

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/server/FilteredBlurServer.java
----------------------------------------------------------------------
diff --git 
a/blur-core/src/main/java/org/apache/blur/server/FilteredBlurServer.java 
b/blur-core/src/main/java/org/apache/blur/server/FilteredBlurServer.java
index 2154916..9318f95 100644
--- a/blur-core/src/main/java/org/apache/blur/server/FilteredBlurServer.java
+++ b/blur-core/src/main/java/org/apache/blur/server/FilteredBlurServer.java
@@ -256,11 +256,6 @@ public class FilteredBlurServer implements Iface {
   }
 
   @Override
-  public Response reconnect(String executionId) throws BlurException, 
TimeoutException, TException {
-    return _iface.reconnect(executionId);
-  }
-
-  @Override
   public void refresh() throws TException {
     _iface.refresh();
   }
@@ -272,18 +267,23 @@ public class FilteredBlurServer implements Iface {
   }
 
   @Override
-  public CommandStatus commandStatus(String executionId) throws BlurException, 
TException {
-    return _iface.commandStatus(executionId);
+  public List<CommandDescriptor> listInstalledCommands() throws BlurException, 
TException {
+    return _iface.listInstalledCommands();
+  }
+
+  @Override
+  public Response reconnect(long instanceExecutionId) throws BlurException, 
TimeoutException, TException {
+    return _iface.reconnect(instanceExecutionId);
   }
 
   @Override
-  public void commandCancel(String executionId) throws BlurException, 
TException {
-    _iface.commandCancel(executionId);
+  public CommandStatus commandStatus(String commandExecutionId) throws 
BlurException, TException {
+    return _iface.commandStatus(commandExecutionId);
   }
 
   @Override
-  public List<CommandDescriptor> listInstalledCommands() throws BlurException, 
TException {
-    return _iface.listInstalledCommands();
+  public void commandCancel(String commandExecutionId) throws BlurException, 
TException {
+    _iface.commandCancel(commandExecutionId);
   }
 
 }

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/thrift/BlurControllerServer.java
----------------------------------------------------------------------
diff --git 
a/blur-core/src/main/java/org/apache/blur/thrift/BlurControllerServer.java 
b/blur-core/src/main/java/org/apache/blur/thrift/BlurControllerServer.java
index 477e0d0..23f719c 100644
--- a/blur-core/src/main/java/org/apache/blur/thrift/BlurControllerServer.java
+++ b/blur-core/src/main/java/org/apache/blur/thrift/BlurControllerServer.java
@@ -53,7 +53,6 @@ import org.apache.blur.command.BlurObject;
 import org.apache.blur.command.BlurObjectSerDe;
 import org.apache.blur.command.CommandUtil;
 import org.apache.blur.command.ControllerCommandManager;
-import org.apache.blur.command.ExecutionId;
 import org.apache.blur.command.Response;
 import org.apache.blur.command.Server;
 import org.apache.blur.command.Shard;
@@ -1525,7 +1524,7 @@ public class BlurControllerServer extends TableAdmin 
implements Iface {
       return CommandUtil.fromObjectToThrift(response, _serDe);
     } catch (Exception e) {
       if (e instanceof org.apache.blur.command.TimeoutException) {
-        throw new TimeoutException(((org.apache.blur.command.TimeoutException) 
e).getExecutionId().getId());
+        throw new TimeoutException(((org.apache.blur.command.TimeoutException) 
e).getInstanceExecutionId());
       }
       LOG.error("Unknown error while trying to execute command [{0}]", e, 
commandName);
       if (e instanceof BlurException) {
@@ -1623,16 +1622,16 @@ public class BlurControllerServer extends TableAdmin 
implements Iface {
   }
 
   @Override
-  public org.apache.blur.thrift.generated.Response reconnect(String 
executionId) throws BlurException,
+  public org.apache.blur.thrift.generated.Response reconnect(long 
instanceExecutionId) throws BlurException,
       TimeoutException, TException {
     try {
-      Response response = _commandManager.reconnect(new 
ExecutionId(executionId));
+      Response response = _commandManager.reconnect(instanceExecutionId);
       return CommandUtil.fromObjectToThrift(response, _serDe);
     } catch (Exception e) {
       if (e instanceof org.apache.blur.command.TimeoutException) {
-        throw new TimeoutException(((org.apache.blur.command.TimeoutException) 
e).getExecutionId().getId());
+        throw new TimeoutException(((org.apache.blur.command.TimeoutException) 
e).getInstanceExecutionId());
       }
-      LOG.error("Unknown error while trying to reconnect to executing command 
[{0}]", e, executionId);
+      LOG.error("Unknown error while trying to reconnect to executing command 
[{0}]", e, instanceExecutionId);
       if (e instanceof BlurException) {
         throw (BlurException) e;
       }
@@ -1665,12 +1664,12 @@ public class BlurControllerServer extends TableAdmin 
implements Iface {
   }
 
   @Override
-  public CommandStatus commandStatus(String executionId) throws BlurException, 
TException {
+  public CommandStatus commandStatus(String commandExecutionId) throws 
BlurException, TException {
     throw new BException("Not Implemented");
   }
 
   @Override
-  public void commandCancel(String executionId) throws BlurException, 
TException {
+  public void commandCancel(String commandExecutionId) throws BlurException, 
TException {
     throw new BException("Not Implemented");
   }
 

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/thrift/BlurShardServer.java
----------------------------------------------------------------------
diff --git 
a/blur-core/src/main/java/org/apache/blur/thrift/BlurShardServer.java 
b/blur-core/src/main/java/org/apache/blur/thrift/BlurShardServer.java
index c1a5b16..135feaf 100644
--- a/blur-core/src/main/java/org/apache/blur/thrift/BlurShardServer.java
+++ b/blur-core/src/main/java/org/apache/blur/thrift/BlurShardServer.java
@@ -36,7 +36,6 @@ import org.apache.blur.command.BlurObject;
 import org.apache.blur.command.BlurObjectSerDe;
 import org.apache.blur.command.CommandStatusStateEnum;
 import org.apache.blur.command.CommandUtil;
-import org.apache.blur.command.ExecutionId;
 import org.apache.blur.command.Response;
 import org.apache.blur.command.ShardCommandManager;
 import org.apache.blur.concurrent.Executors;
@@ -617,7 +616,7 @@ public class BlurShardServer extends TableAdmin implements 
Iface {
       return CommandUtil.fromObjectToThrift(response, _serDe);
     } catch (Exception e) {
       if (e instanceof org.apache.blur.command.TimeoutException) {
-        throw new TimeoutException(((org.apache.blur.command.TimeoutException) 
e).getExecutionId().getId());
+        throw new TimeoutException(((org.apache.blur.command.TimeoutException) 
e).getInstanceExecutionId());
       }
       LOG.error("Unknown error while trying to execute command [{0}] for table 
[{1}]", e, commandName);
       if (e instanceof BlurException) {
@@ -632,16 +631,16 @@ public class BlurShardServer extends TableAdmin 
implements Iface {
   }
 
   @Override
-  public org.apache.blur.thrift.generated.Response reconnect(String 
executionId) throws BlurException,
+  public org.apache.blur.thrift.generated.Response reconnect(long 
instanceExecutionId) throws BlurException,
       TimeoutException, TException {
     try {
-      Response response = _commandManager.reconnect(new 
ExecutionId(executionId));
+      Response response = _commandManager.reconnect(instanceExecutionId);
       return CommandUtil.fromObjectToThrift(response, _serDe);
     } catch (Exception e) {
       if (e instanceof org.apache.blur.command.TimeoutException) {
-        throw new TimeoutException(((org.apache.blur.command.TimeoutException) 
e).getExecutionId().getId());
+        throw new TimeoutException(((org.apache.blur.command.TimeoutException) 
e).getInstanceExecutionId());
       }
-      LOG.error("Unknown error while trying to reconnect to executing command 
[{0}]", e, executionId);
+      LOG.error("Unknown error while trying to reconnect to executing command 
[{0}]", e, instanceExecutionId);
       if (e instanceof BlurException) {
         throw (BlurException) e;
       }
@@ -679,14 +678,14 @@ public class BlurShardServer extends TableAdmin 
implements Iface {
   }
 
   @Override
-  public CommandStatus commandStatus(String executionId) throws BlurException, 
TException {
+  public CommandStatus commandStatus(String commandExecutionId) throws 
BlurException, TException {
     throw new BException("Not Implemented");
   }
 
   @Override
-  public void commandCancel(String executionId) throws BlurException, 
TException {
+  public void commandCancel(String commandExecutionId) throws BlurException, 
TException {
     try {
-      _commandManager.cancel(executionId);
+      _commandManager.cancelCommand(commandExecutionId);
     } catch (Exception e) {
       throw new BException(e.getMessage(), e);
     }
@@ -695,4 +694,5 @@ public class BlurShardServer extends TableAdmin implements 
Iface {
   private CommandStatusStateEnum toCommandStatus(CommandStatusState state) {
     return CommandStatusStateEnum.valueOf(state.name());
   }
+
 }

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/test/java/org/apache/blur/command/ShardCommandManagerTest.java
----------------------------------------------------------------------
diff --git 
a/blur-core/src/test/java/org/apache/blur/command/ShardCommandManagerTest.java 
b/blur-core/src/test/java/org/apache/blur/command/ShardCommandManagerTest.java
index 83dd232..d2b4eed 100644
--- 
a/blur-core/src/test/java/org/apache/blur/command/ShardCommandManagerTest.java
+++ 
b/blur-core/src/test/java/org/apache/blur/command/ShardCommandManagerTest.java
@@ -30,6 +30,7 @@ import java.util.List;
 import java.util.Map;
 import java.util.Set;
 import java.util.SortedSet;
+import java.util.concurrent.CancellationException;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicBoolean;
 
@@ -40,6 +41,8 @@ import org.apache.blur.server.IndexSearcherClosable;
 import org.apache.blur.server.ShardContext;
 import org.apache.blur.server.TableContext;
 import org.apache.blur.server.TableContextFactory;
+import org.apache.blur.thrift.generated.Arguments;
+import org.apache.blur.thrift.generated.BlurException;
 import org.apache.blur.thrift.generated.RowMutation;
 import org.apache.blur.thrift.generated.ShardState;
 import org.apache.blur.thrift.generated.TableDescriptor;
@@ -183,7 +186,7 @@ public class ShardCommandManagerTest {
   @Test
   public void testShardCommandManagerNormalWait() throws IOException, 
TimeoutException, ExceptionCollector {
     Response response;
-    ExecutionId executionId = null;
+    Long instanceExecutionId = null;
 
     BlurObject args = new BlurObject();
     args.put("table", "test");
@@ -197,15 +200,15 @@ public class ShardCommandManagerTest {
         fail();
       }
       try {
-        if (executionId == null) {
+        if (instanceExecutionId == null) {
           TableContextFactory tableContextFactory = getTableContextFactory();
           response = _manager.execute(tableContextFactory, "wait", 
argumentOverlay);
         } else {
-          response = _manager.reconnect(executionId);
+          response = _manager.reconnect(instanceExecutionId);
         }
         break;
       } catch (TimeoutException te) {
-        executionId = te.getExecutionId();
+        instanceExecutionId = te.getInstanceExecutionId();
       }
     }
     System.out.println(response);
@@ -229,24 +232,66 @@ public class ShardCommandManagerTest {
   }
 
   @Test
-  public void testShardCommandManagerNormalWithCancel() throws IOException, 
TimeoutException, ExceptionCollector {
-    Response response;
-    ExecutionId executionId = null;
+  public void testShardCommandManagerNormalWithCancel() throws IOException, 
TimeoutException, ExceptionCollector,
+      BlurException, InterruptedException {
 
-    BlurObject args = new BlurObject();
-    args.put("table", "test");
-    args.put("seconds", 5);
+    String commandExecutionId = "TEST_COMMAND_ID1";
 
-    ArgumentOverlay argumentOverlay = new ArgumentOverlay(args, new 
BlurObjectSerDe());
+    BlurObjectSerDe serDe = new BlurObjectSerDe();
+    WaitForSeconds waitForSeconds = new WaitForSeconds();
+    waitForSeconds.setTable("test");
+    waitForSeconds.setSeconds(5);
+    waitForSeconds.setCommandExecutionId(commandExecutionId);
 
-    try {
-      TableContextFactory tableContextFactory = getTableContextFactory();
-      response = _manager.execute(tableContextFactory, "wait", 
argumentOverlay);
-    } catch (TimeoutException te) {
-      _manager.cancel(te.getExecutionId());
-      // some how validate the threads have cancelled.
-    }
+    Arguments arguments = CommandUtil.toArguments(waitForSeconds, serDe);
+    BlurObject args = CommandUtil.toBlurObject(arguments);
+    System.out.println(args.toString(1));
+    final ArgumentOverlay argumentOverlay = new ArgumentOverlay(args, serDe);
 
+    final AtomicBoolean fail = new AtomicBoolean();
+    final AtomicBoolean running = new AtomicBoolean(true);
+
+    new Thread(new Runnable() {
+      @Override
+      public void run() {
+        TableContextFactory tableContextFactory = getTableContextFactory();
+        Long instanceExecutionId = null;
+        while (true) {
+          try {
+            Response response;
+            if (instanceExecutionId == null) {
+              response = _manager.execute(tableContextFactory, "wait", 
argumentOverlay);
+            } else {
+              response = _manager.reconnect(instanceExecutionId);
+            }
+            fail.set(true);
+            System.out.println(response);
+            return;
+          } catch (IOException e) {
+            if (e.getCause() instanceof CancellationException) {
+              return;  
+            }
+            e.printStackTrace();
+            fail.set(true);
+            return;
+          } catch (TimeoutException e) {
+            instanceExecutionId = e.getInstanceExecutionId();
+          } catch (Exception e) {
+            e.printStackTrace();
+            fail.set(true);
+            return;
+          } finally {
+            running.set(false);
+          }
+        }
+      }
+    }).start();
+    Thread.sleep(1000);
+    _manager.cancelCommand(commandExecutionId);
+    Thread.sleep(5000);
+    if (fail.get() || running.get()) {
+      fail("Fail [" + fail.get() + "] Running [" + running.get() + "]");
+    }
   }
 
   private TableContextFactory getTableContextFactory() {

Reply via email to