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); }
