Repository: airavata
Updated Branches:
  refs/heads/master a665aa56a -> 264bbeb09


fixing the email monitoring to use cachedthreadpool


Project: http://git-wip-us.apache.org/repos/asf/airavata/repo
Commit: http://git-wip-us.apache.org/repos/asf/airavata/commit/264bbeb0
Tree: http://git-wip-us.apache.org/repos/asf/airavata/tree/264bbeb0
Diff: http://git-wip-us.apache.org/repos/asf/airavata/diff/264bbeb0

Branch: refs/heads/master
Commit: 264bbeb09624aa8f41b4f32959ceb08e775b889e
Parents: a665aa5
Author: Lahiru Gunathilake <[email protected]>
Authored: Fri Apr 24 14:51:20 2015 -0400
Committer: Lahiru Gunathilake <[email protected]>
Committed: Fri Apr 24 14:51:20 2015 -0400

----------------------------------------------------------------------
 .../client/samples/CreateLaunchExperiment.java  | 14 ++++++------
 .../tools/RegisterSampleApplications.java       | 11 +++++----
 .../airavata/common/utils/AiravataZKUtils.java  | 20 ++++++++--------
 .../apache/airavata/gfac/server/GfacServer.java |  2 +-
 .../airavata/gfac/server/GfacServerHandler.java |  8 ++-----
 .../airavata/gfac/core/cpi/BetterGfacImpl.java  |  4 ++--
 .../gfac/core/utils/GFacThreadPoolExecutor.java |  2 +-
 .../airavata/gfac/core/utils/GFacUtils.java     | 24 ++++++--------------
 .../gfac/monitor/email/EmailBasedMonitor.java   |  8 +------
 .../monitor/impl/pull/qstat/HPCPullMonitor.java |  8 +++----
 10 files changed, 41 insertions(+), 60 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/airavata/blob/264bbeb0/airavata-api/airavata-client-sdks/java-client-samples/src/main/java/org/apache/airavata/client/samples/CreateLaunchExperiment.java
----------------------------------------------------------------------
diff --git 
a/airavata-api/airavata-client-sdks/java-client-samples/src/main/java/org/apache/airavata/client/samples/CreateLaunchExperiment.java
 
b/airavata-api/airavata-client-sdks/java-client-samples/src/main/java/org/apache/airavata/client/samples/CreateLaunchExperiment.java
index 71ad870..62cd190 100644
--- 
a/airavata-api/airavata-client-sdks/java-client-samples/src/main/java/org/apache/airavata/client/samples/CreateLaunchExperiment.java
+++ 
b/airavata-api/airavata-client-sdks/java-client-samples/src/main/java/org/apache/airavata/client/samples/CreateLaunchExperiment.java
@@ -58,10 +58,10 @@ public class CreateLaunchExperiment {
     private static final String DEFAULT_GATEWAY = "php_reference_gateway";
     private static Airavata.Client airavataClient;
 
-    private static String echoAppId = 
"Echo_4d075102-3591-4ad5-a218-c7e63471ba4a";
+    private static String echoAppId = 
"Echo_418d87fe-c720-4079-acd9-7e78aaa6cf76";
     private static String mpiAppId = 
"HelloMPI_bfd56d58-6085-4b7f-89fc-646576830518";
     private static String wrfAppId = 
"WRF_7ad5da38-c08b-417c-a9ea-da9298839762";
-    private static String amberAppId = 
"Amber_6da3c39f-c435-413a-bb90-f14f87e8d6c3";
+    private static String amberAppId = 
"Amber_3864a9dc-de75-4bd1-8e8e-a0d1c006a35c";
     private static String gromacsAppId = 
"GROMACS_05622038-9edd-4cb1-824e-0b7cb993364b";
     private static String espressoAppId = 
"ESPRESSO_10cc2820-5d0b-4c63-9546-8a8b595593c1";
     private static String lammpsAppId = 
"LAMMPS_2472685b-8acf-497e-aafe-cc66fe5f4cb6";
@@ -166,14 +166,14 @@ public class CreateLaunchExperiment {
 //                final String expId = createExperimentForSSHHost(airavata);
 //                final String expId = 
createEchoExperimentForFSD(airavataClient);
 //                final String expId = 
createMPIExperimentForFSD(airavataClient);
-               final String expId = 
createEchoExperimentForStampede(airavataClient);
+//               final String expId = 
createEchoExperimentForStampede(airavataClient);
 //                final String expId = 
createEchoExperimentForTrestles(airavataClient);
 //                final String expId = 
createExperimentEchoForLocalHost(airavataClient);
 //                final String expId = 
createExperimentWRFTrestles(airavataClient);
 //                final String expId = createExperimentForBR2(airavataClient);
 //                final String expId = 
createExperimentForBR2Amber(airavataClient);
 //                final String expId = 
createExperimentWRFStampede(airavataClient);
-//                final String expId = 
createExperimentForStampedeAmber(airavataClient);
+                final String expId = 
createExperimentForStampedeAmber(airavataClient);
 //                String expId = 
createExperimentForTrestlesAmber(airavataClient);
 //                final String expId = 
createExperimentGROMACSStampede(airavataClient);
 //                final String expId = 
createExperimentESPRESSOStampede(airavataClient);
@@ -1323,11 +1323,11 @@ public class CreateLaunchExperiment {
 //                     }
             for (InputDataObjectType inputDataObjectType : exInputs) {
                 if 
(inputDataObjectType.getName().equalsIgnoreCase("Heat_Restart_File")) {
-                    
inputDataObjectType.setValue("/Users/chathuri/dev/airavata/source/php/inputs/AMBER_FILES/02_Heat.rst");
+                    
inputDataObjectType.setValue("/Users/lginnali/Downloads/02_Heat.rst");
                 } else if 
(inputDataObjectType.getName().equalsIgnoreCase("Production_Control_File")) {
-                    
inputDataObjectType.setValue("/Users/chathuri/dev/airavata/source/php/inputs/AMBER_FILES/03_Prod.in");
+                    
inputDataObjectType.setValue("/Users/lginnali/Downloads/03_Prod.in");
                 } else if 
(inputDataObjectType.getName().equalsIgnoreCase("Parameter_Topology_File")) {
-                    
inputDataObjectType.setValue("/Users/chathuri/dev/airavata/source/php/inputs/AMBER_FILES/prmtop");
+                    
inputDataObjectType.setValue("/Users/lginnali/Downloads/prmtop");
                 }
             }
 

http://git-wip-us.apache.org/repos/asf/airavata/blob/264bbeb0/airavata-api/airavata-client-sdks/java-client-samples/src/main/java/org/apache/airavata/client/tools/RegisterSampleApplications.java
----------------------------------------------------------------------
diff --git 
a/airavata-api/airavata-client-sdks/java-client-samples/src/main/java/org/apache/airavata/client/tools/RegisterSampleApplications.java
 
b/airavata-api/airavata-client-sdks/java-client-samples/src/main/java/org/apache/airavata/client/tools/RegisterSampleApplications.java
index e53074e..3464246 100644
--- 
a/airavata-api/airavata-client-sdks/java-client-samples/src/main/java/org/apache/airavata/client/tools/RegisterSampleApplications.java
+++ 
b/airavata-api/airavata-client-sdks/java-client-samples/src/main/java/org/apache/airavata/client/tools/RegisterSampleApplications.java
@@ -200,7 +200,7 @@ public class RegisterSampleApplications {
             System.out.println("\n #### Registering XSEDE Computational 
Resources #### \n");
 
             //Register Stampede
-            List<BatchQueue> stampedeQueues = new ArrayList<>();
+            List<BatchQueue> stampedeQueues = new ArrayList<BatchQueue>();
             BatchQueue normalQueue = createBatchQueue("normal", "Normal 
Queue", 2880, 256, 4000, 50, 0);
             BatchQueue developmentQueue = createBatchQueue("development", 
"Development Queue", 120, 16, 4000, 1, 0);
             stampedeQueues.add(normalQueue);
@@ -211,7 +211,7 @@ public class RegisterSampleApplications {
             System.out.println("Stampede Resource Id is " + 
stampedeResourceId);
 
             //Register Trestles
-            List<BatchQueue> trestlesQueues = new ArrayList<>();
+            List<BatchQueue> trestlesQueues = new ArrayList<BatchQueue>();
             BatchQueue normalQueue_tr = createBatchQueue("normal", "Normal 
Queue", 2880, 32, 32, 50, 0);
             BatchQueue sharedQueue_tr = createBatchQueue("shared", "Shared 
Queue", 2880, 4, 32, 50, 0);
             trestlesQueues.add(normalQueue_tr);
@@ -221,7 +221,7 @@ public class RegisterSampleApplications {
             System.out.println("Trestles Resource Id is " + 
trestlesResourceId);
 
             //Register BigRedII
-            List<BatchQueue> br2Queues = new ArrayList<>();
+            List<BatchQueue> br2Queues = new ArrayList<BatchQueue>();
             BatchQueue normalQueue_br2 = createBatchQueue("normal", "Normal 
Queue", 2880, 340, 2048, 50, 0);
             BatchQueue serial_br2 = createBatchQueue("serial", "Normal Queue", 
10080, 340, 2048, 50, 0);
             br2Queues.add(normalQueue_br2);
@@ -234,7 +234,7 @@ public class RegisterSampleApplications {
             System.out.println("FSd Resource Id: "+fsdResourceId);
 
             //Register Alamo
-            List<BatchQueue> alamoQueues = new ArrayList<>();
+            List<BatchQueue> alamoQueues = new ArrayList<BatchQueue>();
             alamoResourceId = registerComputeHost("alamo.uthscsa.edu", "Alamo 
Cluster",
                     ResourceJobManagerType.PBS, "push", "/usr/bin/", 
SecurityProtocol.SSH_KEYS, 22, "/usr/bin/mpiexec -np", alamoQueues);
             System.out.println("Alamo Cluster " + alamoResourceId);
@@ -251,7 +251,7 @@ public class RegisterSampleApplications {
             System.out.println("\n #### Registering Non-XSEDE Computational 
Resources #### \n");
 
             //Register LSF resource
-            List<BatchQueue> lsfQueues = new ArrayList<>();
+            List<BatchQueue> lsfQueues = new ArrayList<BatchQueue>();
             lsfResourceId = registerComputeHost("ghpcc06.umassrc.org", "LSF 
Cluster",
                     ResourceJobManagerType.LSF, "push", "source 
/etc/bashrc;/lsf/9.1/linux2.6-glibc2.3-x86_64/bin", SecurityProtocol.SSH_KEYS, 
22, "mpiexec", lsfQueues);
             System.out.println("LSF Resource Id is " + lsfResourceId);
@@ -1379,6 +1379,7 @@ public class RegisterSampleApplications {
         sshJobSubmission.setResourceJobManager(resourceJobManager);
         sshJobSubmission.setSecurityProtocol(securityProtocol);
         sshJobSubmission.setSshPort(portNumber);
+        
sshJobSubmission.setMonitorMode(MonitorMode.JOB_EMAIL_NOTIFICATION_MONITOR);
         airavataClient.addSSHJobSubmissionDetails(computeResourceId, 1, 
sshJobSubmission);
 
         SCPDataMovement scpDataMovement = new SCPDataMovement();

http://git-wip-us.apache.org/repos/asf/airavata/blob/264bbeb0/modules/commons/utils/src/main/java/org/apache/airavata/common/utils/AiravataZKUtils.java
----------------------------------------------------------------------
diff --git 
a/modules/commons/utils/src/main/java/org/apache/airavata/common/utils/AiravataZKUtils.java
 
b/modules/commons/utils/src/main/java/org/apache/airavata/common/utils/AiravataZKUtils.java
index 73d9c27..92364cc 100644
--- 
a/modules/commons/utils/src/main/java/org/apache/airavata/common/utils/AiravataZKUtils.java
+++ 
b/modules/commons/utils/src/main/java/org/apache/airavata/common/utils/AiravataZKUtils.java
@@ -50,14 +50,14 @@ public class AiravataZKUtils implements Watcher {
 
     }
 
-    public static String getExpZnodePath(String experimentId, String taskId) 
throws ApplicationSettingsException {
+    public static String getExpZnodePath(String experimentId) throws 
ApplicationSettingsException {
         return 
ServerSettings.getSetting(Constants.ZOOKEEPER_GFAC_EXPERIMENT_NODE) +
                 File.separator +
                 
ServerSettings.getSetting(Constants.ZOOKEEPER_GFAC_SERVER_NAME) + File.separator
                 + experimentId;
     }
 
-    public static String getExpZnodeHandlerPath(String experimentId, String 
taskId, String className) throws ApplicationSettingsException {
+    public static String getExpZnodeHandlerPath(String experimentId, String 
className) throws ApplicationSettingsException {
         return 
ServerSettings.getSetting(Constants.ZOOKEEPER_GFAC_EXPERIMENT_NODE) +
                 File.separator +
                 
ServerSettings.getSetting(Constants.ZOOKEEPER_GFAC_SERVER_NAME) + File.separator
@@ -73,26 +73,26 @@ public class AiravataZKUtils implements Watcher {
         return 
Integer.parseInt(ServerSettings.getSetting(Constants.ZOOKEEPER_TIMEOUT,"30000"));
     }
 
-    public static String getExpStatePath(String experimentId, String taskId) 
throws ApplicationSettingsException {
-        return AiravataZKUtils.getExpZnodePath(experimentId, taskId) +
+    public static String getExpStatePath(String experimentId) throws 
ApplicationSettingsException {
+        return AiravataZKUtils.getExpZnodePath(experimentId) +
                 File.separator +
                 "state";
     }
 
-    public static String getExpTokenId(ZooKeeper zk, String expId, String tId) 
throws ApplicationSettingsException,
+    public static String getExpTokenId(ZooKeeper zk, String expId) throws 
ApplicationSettingsException,
             KeeperException, InterruptedException {
-        Stat exists = zk.exists(getExpZnodePath(expId, tId), false);
+        Stat exists = zk.exists(getExpZnodePath(expId), false);
         if (exists != null) {
-            return new String(zk.getData(getExpZnodePath(expId, tId), false, 
exists));
+            return new String(zk.getData(getExpZnodePath(expId), false, 
exists));
         }
         return null;
     }
 
-    public static String getExpState(ZooKeeper zk, String expId, String tId) 
throws ApplicationSettingsException,
+    public static String getExpState(ZooKeeper zk, String expId) throws 
ApplicationSettingsException,
             KeeperException, InterruptedException {
-        Stat exists = zk.exists(getExpStatePath(expId, tId), false);
+        Stat exists = zk.exists(getExpStatePath(expId), false);
         if (exists != null) {
-            return new String(zk.getData(getExpStatePath(expId, tId), false, 
exists));
+            return new String(zk.getData(getExpStatePath(expId), false, 
exists));
         }
         return null;
     }

http://git-wip-us.apache.org/repos/asf/airavata/blob/264bbeb0/modules/gfac/airavata-gfac-service/src/main/java/org/apache/airavata/gfac/server/GfacServer.java
----------------------------------------------------------------------
diff --git 
a/modules/gfac/airavata-gfac-service/src/main/java/org/apache/airavata/gfac/server/GfacServer.java
 
b/modules/gfac/airavata-gfac-service/src/main/java/org/apache/airavata/gfac/server/GfacServer.java
index 01c12ad..6689c6b 100644
--- 
a/modules/gfac/airavata-gfac-service/src/main/java/org/apache/airavata/gfac/server/GfacServer.java
+++ 
b/modules/gfac/airavata-gfac-service/src/main/java/org/apache/airavata/gfac/server/GfacServer.java
@@ -110,7 +110,7 @@ public class GfacServer implements IServer{
                        setStatus(IServer.ServerStatus.STOPING);
                        server.stop();
                }
-               GFacThreadPoolExecutor.getThreadPool().shutdownNow();
+               GFacThreadPoolExecutor.getCachedThreadPool().shutdownNow();
 
        }
 

http://git-wip-us.apache.org/repos/asf/airavata/blob/264bbeb0/modules/gfac/airavata-gfac-service/src/main/java/org/apache/airavata/gfac/server/GfacServerHandler.java
----------------------------------------------------------------------
diff --git 
a/modules/gfac/airavata-gfac-service/src/main/java/org/apache/airavata/gfac/server/GfacServerHandler.java
 
b/modules/gfac/airavata-gfac-service/src/main/java/org/apache/airavata/gfac/server/GfacServerHandler.java
index dcb8c54..38caef2 100644
--- 
a/modules/gfac/airavata-gfac-service/src/main/java/org/apache/airavata/gfac/server/GfacServerHandler.java
+++ 
b/modules/gfac/airavata-gfac-service/src/main/java/org/apache/airavata/gfac/server/GfacServerHandler.java
@@ -21,7 +21,6 @@
 package org.apache.airavata.gfac.server;
 
 import com.google.common.eventbus.EventBus;
-import edu.uiuc.ncsa.security.delegation.services.Server;
 import org.airavata.appcatalog.cpi.AppCatalog;
 import org.airavata.appcatalog.cpi.AppCatalogException;
 import org.apache.aiaravata.application.catalog.data.impl.AppCatalogFactory;
@@ -29,7 +28,6 @@ import org.apache.airavata.common.exception.AiravataException;
 import org.apache.airavata.common.exception.ApplicationSettingsException;
 import org.apache.airavata.common.logger.AiravataLogger;
 import org.apache.airavata.common.logger.AiravataLoggerFactory;
-import org.apache.airavata.common.utils.*;
 import org.apache.airavata.common.utils.AiravataZKUtils;
 import org.apache.airavata.common.utils.Constants;
 import org.apache.airavata.common.utils.MonitorPublisher;
@@ -61,8 +59,6 @@ import java.io.File;
 import java.io.IOException;
 import java.util.*;
 import java.util.concurrent.BlockingQueue;
-import java.util.concurrent.ExecutionException;
-import java.util.concurrent.Future;
 import java.util.concurrent.locks.Lock;
 
 
@@ -241,7 +237,7 @@ public class GfacServerHandler implements 
GfacService.Iface, Watcher {
         logger.debugId(experimentId, "Submitted job to the Gfac 
Implementation, experiment {}, task {}, gateway " +
                 "{}", experimentId, taskId, gatewayId);
 
-        GFacThreadPoolExecutor.getThreadPool().execute(inputHandlerWorker);
+        
GFacThreadPoolExecutor.getCachedThreadPool().execute(inputHandlerWorker);
 
         // we immediately return when we have a threadpool
         return true;
@@ -360,7 +356,7 @@ public class GfacServerHandler implements 
GfacService.Iface, Watcher {
 
                     try {
                         
GFacUtils.createExperimentEntryForPassive(event.getExperimentId(), 
event.getTaskId(), zk, experimentNode, nodeName, event.getTokenId(), 
message.getDeliveryTag());
-                        
AiravataZKUtils.getExpStatePath(event.getExperimentId(), event.getTaskId());
+                        
AiravataZKUtils.getExpStatePath(event.getExperimentId());
                         submitJob(event.getExperimentId(), event.getTaskId(), 
event.getGatewayId());
                     } catch (KeeperException e) {
                         logger.error(nodeName + " was interrupted.");

http://git-wip-us.apache.org/repos/asf/airavata/blob/264bbeb0/modules/gfac/gfac-core/src/main/java/org/apache/airavata/gfac/core/cpi/BetterGfacImpl.java
----------------------------------------------------------------------
diff --git 
a/modules/gfac/gfac-core/src/main/java/org/apache/airavata/gfac/core/cpi/BetterGfacImpl.java
 
b/modules/gfac/gfac-core/src/main/java/org/apache/airavata/gfac/core/cpi/BetterGfacImpl.java
index a6df8d5..45e320d 100644
--- 
a/modules/gfac/gfac-core/src/main/java/org/apache/airavata/gfac/core/cpi/BetterGfacImpl.java
+++ 
b/modules/gfac/gfac-core/src/main/java/org/apache/airavata/gfac/core/cpi/BetterGfacImpl.java
@@ -306,7 +306,7 @@ public class BetterGfacImpl implements GFac,Watcher {
         jobExecutionContext.setProperty(Constants.PROP_TOPIC, experimentID);
         jobExecutionContext.setGfac(this);
         jobExecutionContext.setZk(zk);
-        
jobExecutionContext.setCredentialStoreToken(AiravataZKUtils.getExpTokenId(zk, 
experimentID, taskID));
+        
jobExecutionContext.setCredentialStoreToken(AiravataZKUtils.getExpTokenId(zk, 
experimentID));
 
         // handle job submission protocol
         List<JobSubmissionInterface> jobSubmissionInterfaces = 
computeResource.getJobSubmissionInterfaces();
@@ -500,7 +500,7 @@ public class BetterGfacImpl implements GFac,Watcher {
             } else if (stateVal >= 8) {
                 log.info("There is nothing to recover in this job so we do not 
re-submit");
                 ZKUtil.deleteRecursive(zk,
-                        
AiravataZKUtils.getExpZnodePath(jobExecutionContext.getExperimentID(), 
jobExecutionContext.getTaskData().getTaskID()));
+                        
AiravataZKUtils.getExpZnodePath(jobExecutionContext.getExperimentID()));
             } else {
                 // Now we know this is an old Job, so we have to handle things 
gracefully
                 log.info("Re-launching the job in GFac because this is 
re-submitted to GFac");

http://git-wip-us.apache.org/repos/asf/airavata/blob/264bbeb0/modules/gfac/gfac-core/src/main/java/org/apache/airavata/gfac/core/utils/GFacThreadPoolExecutor.java
----------------------------------------------------------------------
diff --git 
a/modules/gfac/gfac-core/src/main/java/org/apache/airavata/gfac/core/utils/GFacThreadPoolExecutor.java
 
b/modules/gfac/gfac-core/src/main/java/org/apache/airavata/gfac/core/utils/GFacThreadPoolExecutor.java
index 3c8d56a..3ce7a50 100644
--- 
a/modules/gfac/gfac-core/src/main/java/org/apache/airavata/gfac/core/utils/GFacThreadPoolExecutor.java
+++ 
b/modules/gfac/gfac-core/src/main/java/org/apache/airavata/gfac/core/utils/GFacThreadPoolExecutor.java
@@ -34,7 +34,7 @@ public class GFacThreadPoolExecutor {
 
     private static ExecutorService threadPool;
 
-    public static ExecutorService getThreadPool() {
+    public static ExecutorService getCachedThreadPool() {
         if(threadPool ==null){
             threadPool = Executors.newCachedThreadPool();
         }

http://git-wip-us.apache.org/repos/asf/airavata/blob/264bbeb0/modules/gfac/gfac-core/src/main/java/org/apache/airavata/gfac/core/utils/GFacUtils.java
----------------------------------------------------------------------
diff --git 
a/modules/gfac/gfac-core/src/main/java/org/apache/airavata/gfac/core/utils/GFacUtils.java
 
b/modules/gfac/gfac-core/src/main/java/org/apache/airavata/gfac/core/utils/GFacUtils.java
index 054e0c3..206bbdd 100644
--- 
a/modules/gfac/gfac-core/src/main/java/org/apache/airavata/gfac/core/utils/GFacUtils.java
+++ 
b/modules/gfac/gfac-core/src/main/java/org/apache/airavata/gfac/core/utils/GFacUtils.java
@@ -20,7 +20,6 @@
  */
 package org.apache.airavata.gfac.core.utils;
 
-import edu.uiuc.ncsa.security.delegation.services.Server;
 import org.airavata.appcatalog.cpi.AppCatalog;
 import org.airavata.appcatalog.cpi.AppCatalogException;
 import org.apache.aiaravata.application.catalog.data.impl.AppCatalogFactory;
@@ -42,7 +41,6 @@ import 
org.apache.airavata.model.appcatalog.appinterface.InputDataObjectType;
 import org.apache.airavata.model.appcatalog.appinterface.OutputDataObjectType;
 import org.apache.airavata.model.appcatalog.computeresource.*;
 import org.apache.airavata.model.workspace.experiment.*;
-import org.apache.airavata.persistance.registry.jpa.impl.RegistryFactory;
 import org.apache.airavata.registry.cpi.ChildDataType;
 import org.apache.airavata.registry.cpi.CompositeIdentifier;
 import org.apache.airavata.registry.cpi.Registry;
@@ -901,8 +899,7 @@ public class GFacUtils {
                        throws ApplicationSettingsException, KeeperException,
                        InterruptedException {
                String expState = AiravataZKUtils.getExpState(zk, 
jobExecutionContext
-                               .getExperimentID(), 
jobExecutionContext.getTaskData()
-                               .getTaskID());
+                               .getExperimentID());
                return 
GfacExperimentState.findByValue(Integer.parseInt(expState));
        }
 
@@ -911,8 +908,7 @@ public class GFacUtils {
                        throws ApplicationSettingsException, KeeperException,
                        InterruptedException {
                String expState = AiravataZKUtils.getExpState(zk, 
jobExecutionContext
-                               .getExperimentID(), 
jobExecutionContext.getTaskData()
-                               .getTaskID());
+                               .getExperimentID());
                if (expState == null) {
                        return -1;
                }
@@ -933,8 +929,7 @@ public class GFacUtils {
                        throws ApplicationSettingsException, KeeperException,
                        InterruptedException {
                String expState = AiravataZKUtils.getExpZnodeHandlerPath(
-                               jobExecutionContext.getExperimentID(), 
jobExecutionContext
-                                               .getTaskData().getTaskID(), 
className);
+                               jobExecutionContext.getExperimentID(), 
className);
                Stat exists = zk.exists(expState, false);
                if (exists == null) {
                        zk.create(expState, new byte[0], 
ZooDefs.Ids.OPEN_ACL_UNSAFE,
@@ -970,8 +965,7 @@ public class GFacUtils {
                        GfacPluginState state) throws 
ApplicationSettingsException,
                        KeeperException, InterruptedException {
                String expState = AiravataZKUtils.getExpZnodeHandlerPath(
-                               jobExecutionContext.getExperimentID(), 
jobExecutionContext
-                                               .getTaskData().getTaskID(), 
className);
+                               jobExecutionContext.getExperimentID(), 
className);
                Stat exists = zk.exists(expState, false);
                if (exists == null) {
                        zk.create(expState, new byte[0], 
ZooDefs.Ids.OPEN_ACL_UNSAFE,
@@ -1008,8 +1002,7 @@ public class GFacUtils {
                        KeeperException, InterruptedException {
                if(zk.getState().isConnected()) {
                        String expState = 
AiravataZKUtils.getExpZnodeHandlerPath(
-                                       jobExecutionContext.getExperimentID(), 
jobExecutionContext
-                                                       
.getTaskData().getTaskID(), className);
+                                       jobExecutionContext.getExperimentID(), 
className);
 
                        Stat exists = zk.exists(expState + File.separator
                                        + 
AiravataZKUtils.ZK_EXPERIMENT_STATE_NODE, false);
@@ -1030,8 +1023,7 @@ public class GFacUtils {
                        JobExecutionContext jobExecutionContext, String 
className) {
                try {
                        String expState = 
AiravataZKUtils.getExpZnodeHandlerPath(
-                                       jobExecutionContext.getExperimentID(), 
jobExecutionContext
-                                                       
.getTaskData().getTaskID(), className);
+                                       jobExecutionContext.getExperimentID(), 
className);
 
                        Stat exists = zk.exists(expState + File.separator
                                        + 
AiravataZKUtils.ZK_EXPERIMENT_STATE_NODE, false);
@@ -1176,7 +1168,7 @@ public class GFacUtils {
                                                expParent.getVersion());
                        }
 
-            String token = AiravataZKUtils.getExpTokenId(zk, experimentID, 
taskID);
+            String token = AiravataZKUtils.getExpTokenId(zk, experimentID);
             String s = zk.create(newExpNode + File.separator + "state", String
                                                        
.valueOf(GfacExperimentState.LAUNCHED.getValue())
                                                        .getBytes(), 
ZooDefs.Ids.OPEN_ACL_UNSAFE,
@@ -1333,7 +1325,6 @@ public class GFacUtils {
                                String expZnodeHandlerPath = AiravataZKUtils
                                                .getExpZnodeHandlerPath(
                                                                
jobExecutionContext.getExperimentID(),
-                                                               
jobExecutionContext.getTaskData().getTaskID(),
                                                                className);
                                Stat exists = zk.exists(expZnodeHandlerPath, 
false);
                                zk.setData(expZnodeHandlerPath, 
data.toString().getBytes(),
@@ -1352,7 +1343,6 @@ public class GFacUtils {
                        String expZnodeHandlerPath = AiravataZKUtils
                                        .getExpZnodeHandlerPath(
                                                        
jobExecutionContext.getExperimentID(),
-                                                       
jobExecutionContext.getTaskData().getTaskID(),
                                                        className);
                        Stat exists = zk.exists(expZnodeHandlerPath, false);
                        return new String(jobExecutionContext.getZk().getData(

http://git-wip-us.apache.org/repos/asf/airavata/blob/264bbeb0/modules/gfac/gfac-monitor/gfac-email-monitor/src/main/java/org/apache/airavata/gfac/monitor/email/EmailBasedMonitor.java
----------------------------------------------------------------------
diff --git 
a/modules/gfac/gfac-monitor/gfac-email-monitor/src/main/java/org/apache/airavata/gfac/monitor/email/EmailBasedMonitor.java
 
b/modules/gfac/gfac-monitor/gfac-email-monitor/src/main/java/org/apache/airavata/gfac/monitor/email/EmailBasedMonitor.java
index 5fce761..393b3cf 100644
--- 
a/modules/gfac/gfac-monitor/gfac-email-monitor/src/main/java/org/apache/airavata/gfac/monitor/email/EmailBasedMonitor.java
+++ 
b/modules/gfac/gfac-monitor/gfac-email-monitor/src/main/java/org/apache/airavata/gfac/monitor/email/EmailBasedMonitor.java
@@ -20,9 +20,7 @@
 */
 package org.apache.airavata.gfac.monitor.email;
 
-import org.apache.aiaravata.application.catalog.data.model.ResourceJobManager;
 import org.apache.airavata.common.exception.AiravataException;
-import org.apache.airavata.common.exception.ApplicationSettingsException;
 import org.apache.airavata.common.logger.AiravataLogger;
 import org.apache.airavata.common.logger.AiravataLoggerFactory;
 import org.apache.airavata.common.utils.ServerSettings;
@@ -260,11 +258,7 @@ public class EmailBasedMonitor implements Runnable{
         }
 
         if (runOutHandlers) {
-            try {
-                GFacThreadPoolExecutor.getFixedThreadPool().submit(new 
OutHandlerWorker(jEC, BetterGfacImpl.getMonitorPublisher()));
-            } catch (ApplicationSettingsException e) {
-                log.error(e.getMessage(), e);
-            }
+                GFacThreadPoolExecutor.getCachedThreadPool().execute(new 
OutHandlerWorker(jEC, BetterGfacImpl.getMonitorPublisher()));
         }
         publishJobStatusChange(jEC);
     }

http://git-wip-us.apache.org/repos/asf/airavata/blob/264bbeb0/modules/gfac/gfac-monitor/gfac-hpc-monitor/src/main/java/org/apache/airavata/gfac/monitor/impl/pull/qstat/HPCPullMonitor.java
----------------------------------------------------------------------
diff --git 
a/modules/gfac/gfac-monitor/gfac-hpc-monitor/src/main/java/org/apache/airavata/gfac/monitor/impl/pull/qstat/HPCPullMonitor.java
 
b/modules/gfac/gfac-monitor/gfac-hpc-monitor/src/main/java/org/apache/airavata/gfac/monitor/impl/pull/qstat/HPCPullMonitor.java
index 04f9e7d..3442367 100644
--- 
a/modules/gfac/gfac-monitor/gfac-hpc-monitor/src/main/java/org/apache/airavata/gfac/monitor/impl/pull/qstat/HPCPullMonitor.java
+++ 
b/modules/gfac/gfac-monitor/gfac-hpc-monitor/src/main/java/org/apache/airavata/gfac/monitor/impl/pull/qstat/HPCPullMonitor.java
@@ -197,7 +197,7 @@ public class HPCPullMonitor extends PullMonitor {
                                 sendNotification(iMonitorID);
                                 logger.info("To avoid timing issues we sleep 
sometime and try to retrieve output files");
                                 Thread.sleep(10000);
-                                
GFacThreadPoolExecutor.getThreadPool().execute(new OutHandlerWorker(gfac, 
iMonitorID, publisher));
+                                
GFacThreadPoolExecutor.getCachedThreadPool().execute(new OutHandlerWorker(gfac, 
iMonitorID, publisher));
                                 break;
                             }
                         }
@@ -225,7 +225,7 @@ public class HPCPullMonitor extends PullMonitor {
                                     sendNotification(iMonitorID);
                                     logger.info("To avoid timing issues we 
sleep sometime and try to retrieve output files");
                                     Thread.sleep(10000);
-                                    
GFacThreadPoolExecutor.getThreadPool().execute(new OutHandlerWorker(gfac, 
iMonitorID, publisher));
+                                    
GFacThreadPoolExecutor.getCachedThreadPool().execute(new OutHandlerWorker(gfac, 
iMonitorID, publisher));
                                     break;
                                 }
                             }
@@ -250,7 +250,7 @@ public class HPCPullMonitor extends PullMonitor {
                             removeList.add(iMonitorID);
                             logger.info("PULL Notification is complete: 
marking the Job as ************COMPLETE************ experiment {}, task {}, job 
name {} .",
                                     iMonitorID.getExperimentID(), 
iMonitorID.getTaskID(), iMonitorID.getJobName());
-                            GFacThreadPoolExecutor.getThreadPool().execute(new 
OutHandlerWorker(gfac, iMonitorID, publisher));
+                            
GFacThreadPoolExecutor.getCachedThreadPool().execute(new OutHandlerWorker(gfac, 
iMonitorID, publisher));
                         }
                         
iMonitorID.setStatus(jobStatuses.get(iMonitorID.getJobID() + "," + 
iMonitorID.getJobName()));    //IMPORTANT this is not a simple setter we have a 
logic
                         iMonitorID.setLastMonitored(new Timestamp((new 
Date()).getTime()));
@@ -288,7 +288,7 @@ public class HPCPullMonitor extends PullMonitor {
                                 sendNotification(iMonitorID);
 //                                CommonUtils.removeMonitorFromQueue(take, 
iMonitorID);
                                 removeList.add(iMonitorID);
-                                
GFacThreadPoolExecutor.getThreadPool().execute(new OutHandlerWorker(gfac, 
iMonitorID, publisher));
+                                
GFacThreadPoolExecutor.getCachedThreadPool().execute(new OutHandlerWorker(gfac, 
iMonitorID, publisher));
                             } else {
                                 iMonitorID.setFailedCount(0);
                             }

Reply via email to