Repository: airavata Updated Branches: refs/heads/master d35226d88 -> 5b118f8c4
modified provider sources to use appcatalog types Project: http://git-wip-us.apache.org/repos/asf/airavata/repo Commit: http://git-wip-us.apache.org/repos/asf/airavata/commit/28a235e8 Tree: http://git-wip-us.apache.org/repos/asf/airavata/tree/28a235e8 Diff: http://git-wip-us.apache.org/repos/asf/airavata/diff/28a235e8 Branch: refs/heads/master Commit: 28a235e87c5820ab2d5e7e9849be5f041823b35d Parents: 7f4faeb Author: msmemon <[email protected]> Authored: Wed Dec 3 13:24:55 2014 +0100 Committer: msmemon <[email protected]> Committed: Wed Dec 3 13:24:55 2014 +0100 ---------------------------------------------------------------------- .../gfac/bes/handlers/AbstractSMSHandler.java | 2 - .../gfac/bes/provider/impl/BESProvider.java | 81 +++++---- .../gfac/bes/utils/ApplicationProcessor.java | 60 ++++--- .../gfac/bes/utils/DataServiceInfo.java | 62 ------- .../gfac/bes/utils/DataTransferrer.java | 158 +++++----------- .../airavata/gfac/bes/utils/JSDLGenerator.java | 84 ++------- .../gfac/bes/utils/ResourceProcessor.java | 149 +++++---------- .../gfac/bes/utils/UASDataStagingProcessor.java | 180 ++++++------------- 8 files changed, 230 insertions(+), 546 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/airavata/blob/28a235e8/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/handlers/AbstractSMSHandler.java ---------------------------------------------------------------------- diff --git a/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/handlers/AbstractSMSHandler.java b/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/handlers/AbstractSMSHandler.java index 71ca0db..a23a096 100644 --- a/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/handlers/AbstractSMSHandler.java +++ b/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/handlers/AbstractSMSHandler.java @@ -17,8 +17,6 @@ import org.apache.airavata.gfac.core.utils.GFacUtils; import org.apache.airavata.model.appcatalog.computeresource.*; import org.apache.airavata.model.workspace.experiment.CorrectiveAction; import org.apache.airavata.model.workspace.experiment.ErrorCategory; -import org.apache.airavata.schemas.gfac.JobDirectoryModeDocument.JobDirectoryMode; -import org.apache.airavata.schemas.gfac.UnicoreHostType; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.w3.x2005.x08.addressing.EndpointReferenceType; http://git-wip-us.apache.org/repos/asf/airavata/blob/28a235e8/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/provider/impl/BESProvider.java ---------------------------------------------------------------------- diff --git a/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/provider/impl/BESProvider.java b/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/provider/impl/BESProvider.java index 7cf2d7c..6fdadfb 100644 --- a/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/provider/impl/BESProvider.java +++ b/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/provider/impl/BESProvider.java @@ -115,8 +115,7 @@ public class BESProvider extends AbstractProvider implements GFacProvider, UnicoreJobSubmission unicoreJobSubmission = GFacUtils.getUnicoreJobSubmission(interfaceId); factoryUrl = unicoreJobSubmission.getUnicoreEndPointURL(); } - EndpointReferenceType eprt = EndpointReferenceType.Factory - .newInstance(); + EndpointReferenceType eprt = EndpointReferenceType.Factory.newInstance(); eprt.addNewAddress().setStringValue(factoryUrl); String userDN = getUserName(jobExecutionContext); @@ -124,20 +123,15 @@ public class BESProvider extends AbstractProvider implements GFacProvider, if (userDN == null || userDN.equalsIgnoreCase("admin")) { userDN = "CN=zdv575, O=Ultrascan Gateway, C=DE"; } - CreateActivityDocument cad = CreateActivityDocument.Factory - .newInstance(); - JobDefinitionDocument jobDefDoc = JobDefinitionDocument.Factory - .newInstance(); + CreateActivityDocument cad = CreateActivityDocument.Factory.newInstance(); + JobDefinitionDocument jobDefDoc = JobDefinitionDocument.Factory.newInstance(); // create storage - StorageCreator storageCreator = new StorageCreator(secProperties, - factoryUrl, 5, null); + StorageCreator storageCreator = new StorageCreator(secProperties, factoryUrl, 5, null); sc = storageCreator.createStorage(); - JobDefinitionType jobDefinition = JSDLGenerator.buildJSDLInstance( - jobExecutionContext, sc.getUrl()).getJobDefinition(); - cad.addNewCreateActivity().addNewActivityDocument() - .setJobDefinition(jobDefinition); + JobDefinitionType jobDefinition = JSDLGenerator.buildJSDLInstance(jobExecutionContext, sc.getUrl()).getJobDefinition(); + cad.addNewCreateActivity().addNewActivityDocument().setJobDefinition(jobDefinition); log.info("JSDL" + jobDefDoc.toString()); // upload files if any @@ -180,28 +174,7 @@ public class BESProvider extends AbstractProvider implements GFacProvider, .getStringValue(), factory.getActivityStatus(activityEpr) .toString())); - // TODO publish the status messages to the message bus - while ((factory.getActivityStatus(activityEpr) != ActivityStateEnumeration.FINISHED) - && (factory.getActivityStatus(activityEpr) != ActivityStateEnumeration.FAILED) - && (factory.getActivityStatus(activityEpr) != ActivityStateEnumeration.CANCELLED)) { - - ActivityStatusType activityStatus = getStatus(factory, activityEpr); - JobState applicationJobStatus = getApplicationJobStatus(activityStatus); - String jobStatusMessage = "Status of job " + jobId + "is " - + applicationJobStatus; - GFacUtils.updateJobStatus(jobExecutionContext, jobDetails, applicationJobStatus); - - jobExecutionContext.getNotifier().publish( - new StatusChangeEvent(jobStatusMessage)); - - // GFacUtils.updateApplicationJobStatus(jobExecutionContext,jobId, - // applicationJobStatus); - try { - Thread.sleep(5000); - } catch (InterruptedException e) { - } - continue; - } + waitUntilDone(factory, activityEpr, jobDetails); ActivityStatusType activityStatus = null; activityStatus = getStatus(factory, activityEpr); @@ -217,11 +190,12 @@ public class BESProvider extends AbstractProvider implements GFacProvider, + activityStatus.getFault().getFaultstring() + "\n EXITCODE: " + activityStatus.getExitCode(); log.info(error); - try { - Thread.sleep(5000); - } catch (InterruptedException e) { - } + + try {Thread.sleep(5000);} catch (InterruptedException e) {} + + //What if job is failed before execution and there are not stdouts generated yet? dt.downloadStdOuts(); + } else if (activityStatus.getState() == ActivityStateEnumeration.CANCELLED) { JobState applicationJobStatus = JobState.CANCELED; String jobStatusMessage = "Status of job " + jobId + "is " @@ -357,8 +331,7 @@ public class BESProvider extends AbstractProvider implements GFacProvider, // } } - protected ActivityStatusType getStatus(FactoryClient fc, - EndpointReferenceType activityEpr) + protected ActivityStatusType getStatus(FactoryClient fc, EndpointReferenceType activityEpr) throws UnknownActivityIdentifierFault { GetActivityStatusesDocument stats = GetActivityStatusesDocument.Factory @@ -427,6 +400,32 @@ public class BESProvider extends AbstractProvider implements GFacProvider, public void cancelJob(JobExecutionContext jobExecutionContext) throws GFacProviderException, GFacException { // TODO Auto-generated method stub - + } + + protected void waitUntilDone(FactoryClient factory, EndpointReferenceType activityEpr, JobDetails jobDetails) throws Exception { + + try { + while ((factory.getActivityStatus(activityEpr) != ActivityStateEnumeration.FINISHED) + && (factory.getActivityStatus(activityEpr) != ActivityStateEnumeration.FAILED) + && (factory.getActivityStatus(activityEpr) != ActivityStateEnumeration.CANCELLED)) { + + ActivityStatusType activityStatus = getStatus(factory, activityEpr); + JobState applicationJobStatus = getApplicationJobStatus(activityStatus); + String jobStatusMessage = "Status of job " + jobId + "is " + applicationJobStatus; + GFacUtils.updateJobStatus(jobExecutionContext, jobDetails, applicationJobStatus); + + jobExecutionContext.getNotifier().publish(new StatusChangeEvent(jobStatusMessage)); + + // GFacUtils.updateApplicationJobStatus(jobExecutionContext,jobId, + // applicationJobStatus); + try { + Thread.sleep(5000); + } catch (InterruptedException e) {} + continue; + } + } catch(Exception e) { + log.error("Error monitoring job status.."); + throw e; + } } } \ No newline at end of file http://git-wip-us.apache.org/repos/asf/airavata/blob/28a235e8/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/ApplicationProcessor.java ---------------------------------------------------------------------- diff --git a/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/ApplicationProcessor.java b/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/ApplicationProcessor.java index ee58565..91c27f9 100644 --- a/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/ApplicationProcessor.java +++ b/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/ApplicationProcessor.java @@ -24,8 +24,6 @@ package org.apache.airavata.gfac.bes.utils; import org.apache.airavata.gfac.core.context.JobExecutionContext; import org.apache.airavata.model.appcatalog.appdeployment.ApplicationDeploymentDescription; import org.apache.airavata.model.appcatalog.appdeployment.ApplicationParallelismType; -import org.apache.airavata.schemas.gfac.ExtendedKeyValueType; -import org.apache.airavata.schemas.gfac.HpcApplicationDeploymentType; import org.ggf.schemas.jsdl.x2005.x11.jsdl.ApplicationType; import org.ggf.schemas.jsdl.x2005.x11.jsdl.JobDefinitionType; import org.ggf.schemas.jsdl.x2005.x11.jsdlPosix.FileNameType; @@ -47,10 +45,9 @@ public class ApplicationProcessor { ApplicationDeploymentDescription appDep= context.getApplicationContext().getApplicationDeploymentDescription(); String appname = context.getApplicationContext().getApplicationInterfaceDescription().getApplicationName(); ApplicationParallelismType parallelism = appDep.getParallelism(); - ApplicationType appType = JSDLUtils.getOrCreateApplication(value); appType.setApplicationName(appname); - JSDLUtils.getOrCreateJobIdentification(value).setJobName(appname); + // if (appDep.getSetEnvironment().size() > 0) { // createApplicationEnvironment(value, appDep.getSetEnvironment(), parallelism); @@ -58,10 +55,11 @@ public class ApplicationProcessor { // String stdout = context.getStandardOutput(); String stderr = context.getStandardError(); + if (appDep.getExecutablePath() != null) { FileNameType fNameType = FileNameType.Factory.newInstance(); fNameType.setStringValue(appDep.getExecutablePath()); - if(parallelism.equals(ApplicationParallelismType.MPI) || parallelism.equals(ApplicationParallelismType.OPENMP_MPI)) { + if(isParallelJob(context)) { JSDLUtils.getOrCreateSPMDApplication(value).setExecutable(fNameType); if (parallelism.equals(ApplicationParallelismType.OPENMP_MPI)){ JSDLUtils.getSPMDApplication(value).setSPMDVariation(SPMDVariations.OpenMPI.value()); @@ -85,11 +83,11 @@ public class ApplicationProcessor { } int totalThreadCount = context.getTaskData().getTaskScheduling().getNumberOfThreads(); + if(totalThreadCount > 0){ ThreadsPerProcessType tpp = ThreadsPerProcessType.Factory.newInstance(); tpp.setStringValue(String.valueOf(totalThreadCount)); JSDLUtils.getSPMDApplication(value).setThreadsPerProcess(tpp); - } if(userName != null) { @@ -134,45 +132,49 @@ public class ApplicationProcessor { public static String getUserNameFromContext(JobExecutionContext jobContext) { if(jobContext.getTaskData() == null) return null; - //FIXME: Discuss to get user and change this + //TODO: Extend unicore model to specify optional unix user id (allocation account) return "admin"; } - public static void addApplicationArgument(JobDefinitionType value, HpcApplicationDeploymentType appDepType, String stringPrm) { - if(isParallelJob(appDepType)) - JSDLUtils.getOrCreateSPMDApplication(value) - .addNewArgument().setStringValue(stringPrm); - else - JSDLUtils.getOrCreatePOSIXApplication(value) - .addNewArgument().setStringValue(stringPrm); - + public static void addApplicationArgument(JobDefinitionType value, JobExecutionContext context, String stringPrm) { + if(isParallelJob(context)){ + JSDLUtils.getOrCreateSPMDApplication(value).addNewArgument().setStringValue(stringPrm); + } + else { + JSDLUtils.getOrCreatePOSIXApplication(value).addNewArgument().setStringValue(stringPrm); + } } - public static String getApplicationStdOut(JobDefinitionType value, HpcApplicationDeploymentType appDepType) throws RuntimeException { - if (isParallelJob(appDepType)) return JSDLUtils.getOrCreateSPMDApplication(value).getOutput().getStringValue(); + public static String getApplicationStdOut(JobDefinitionType value, JobExecutionContext context) throws RuntimeException { + if (isParallelJob(context)) return JSDLUtils.getOrCreateSPMDApplication(value).getOutput().getStringValue(); else return JSDLUtils.getOrCreatePOSIXApplication(value).getOutput().getStringValue(); } - public static String getApplicationStdErr(JobDefinitionType value, HpcApplicationDeploymentType appDepType) throws RuntimeException { - if (isParallelJob(appDepType)) return JSDLUtils.getOrCreateSPMDApplication(value).getError().getStringValue(); + public static String getApplicationStdErr(JobDefinitionType value, JobExecutionContext context) throws RuntimeException { + if (isParallelJob(context)) return JSDLUtils.getOrCreateSPMDApplication(value).getError().getStringValue(); else return JSDLUtils.getOrCreatePOSIXApplication(value).getError().getStringValue(); } public static void createGenericApplication(JobDefinitionType value, String appName) { ApplicationType appType = JSDLUtils.getOrCreateApplication(value); appType.setApplicationName(appName); - JSDLUtils.getOrCreateJobIdentification(value).setJobName(appName); } - - - public static String getValueFromMap(HpcApplicationDeploymentType appDepType, String name) { - ExtendedKeyValueType[] extended = appDepType.getKeyValuePairsArray(); - for(ExtendedKeyValueType e: extended) { - if(e.getName().equalsIgnoreCase(name)) { - return e.getStringValue(); - } + + public static boolean isParallelJob(JobExecutionContext context) { + + ApplicationDeploymentDescription appDep = context.getApplicationContext().getApplicationDeploymentDescription(); + ApplicationParallelismType parallelism = appDep.getParallelism(); + + boolean isParallel = false; + + if(parallelism.equals(ApplicationParallelismType.MPI) || + parallelism.equals(ApplicationParallelismType.OPENMP_MPI) || + parallelism.equals(ApplicationParallelismType.OPENMP )) { + isParallel = true; } - return null; + + return isParallel; } + } http://git-wip-us.apache.org/repos/asf/airavata/blob/28a235e8/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/DataServiceInfo.java ---------------------------------------------------------------------- diff --git a/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/DataServiceInfo.java b/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/DataServiceInfo.java deleted file mode 100644 index b63dcb2..0000000 --- a/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/DataServiceInfo.java +++ /dev/null @@ -1,62 +0,0 @@ -package org.apache.airavata.gfac.bes.utils; - -import java.io.Serializable; - -import org.apache.airavata.gfac.core.context.JobExecutionContext; -import org.apache.airavata.schemas.gfac.JobDirectoryModeDocument.JobDirectoryMode; -import org.apache.airavata.schemas.gfac.UnicoreHostType; -import org.w3.x2005.x08.addressing.EndpointReferenceType; - - -/** - * A value object carrying information about data service access mode. - * */ -public class DataServiceInfo implements BESConstants, Serializable { - - private static final long serialVersionUID = 1L; - - public enum DirectoryAccessMode { - GridFTP, SMSBYTEIO, RNSBYTEIO - } - - /* - * basically only uses information to hold gridftp address or an optional - * pointer to a remote StorageManagementService instance. - */ - private String dataServiceUrl; - - private DirectoryAccessMode directoryAccesMode = DirectoryAccessMode.SMSBYTEIO; - - public DataServiceInfo(JobExecutionContext c) { - JobDirectoryMode.Enum directoryAccess = ((UnicoreHostType)c.getApplicationContext().getHostDescription().getType()).getJobDirectoryMode(); - - switch(directoryAccess.intValue()) { - case JobDirectoryMode.INT_SMS_BYTE_IO: - directoryAccesMode = DirectoryAccessMode.SMSBYTEIO; - EndpointReferenceType s = (EndpointReferenceType) c - .getProperty(PROP_SMS_EPR); - dataServiceUrl = s.getAddress().getStringValue(); - break; - case JobDirectoryMode.INT_GRID_FTP: - case JobDirectoryMode.INT_RNS_BYTE_IO: - default: - directoryAccesMode = DirectoryAccessMode.GridFTP; - break; - } - - } - - public String getDataServiceUrl() { - return dataServiceUrl; - } - - public void setDataServiceUrl(String dataServiceUrl) { - this.dataServiceUrl = dataServiceUrl; - } - - public DirectoryAccessMode getDirectoryAccesMode() { - return directoryAccesMode; - } - - -} \ No newline at end of file http://git-wip-us.apache.org/repos/asf/airavata/blob/28a235e8/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/DataTransferrer.java ---------------------------------------------------------------------- diff --git a/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/DataTransferrer.java b/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/DataTransferrer.java index fa9eb83..4aa6cc1 100644 --- a/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/DataTransferrer.java +++ b/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/DataTransferrer.java @@ -27,22 +27,16 @@ import java.io.FileNotFoundException; import java.io.FileReader; import java.io.IOException; import java.util.ArrayList; -import java.util.HashMap; import java.util.List; -import java.util.Map; -import org.apache.airavata.commons.gfac.type.ActualParameter; -import org.apache.airavata.commons.gfac.type.ApplicationDescription; import org.apache.airavata.gfac.Constants; import org.apache.airavata.gfac.core.context.JobExecutionContext; import org.apache.airavata.gfac.core.provider.GFacProviderException; +import org.apache.airavata.model.appcatalog.appdeployment.ApplicationDeploymentDescription; +import org.apache.airavata.model.appcatalog.appinterface.DataType; +import org.apache.airavata.model.appcatalog.appinterface.InputDataObjectType; +import org.apache.airavata.model.appcatalog.appinterface.OutputDataObjectType; import org.apache.airavata.model.workspace.experiment.TaskDetails; -import org.apache.airavata.schemas.gfac.ApplicationDeploymentDescriptionType; -import org.apache.airavata.schemas.gfac.HpcApplicationDeploymentType; -import org.apache.airavata.schemas.gfac.StringArrayType; -import org.apache.airavata.schemas.gfac.StringParameterType; -import org.apache.airavata.schemas.gfac.URIArrayType; -import org.apache.airavata.schemas.gfac.URIParameterType; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -63,15 +57,9 @@ public class DataTransferrer { public void uploadLocalFiles() throws GFacProviderException { - Map<String, Object> inputParams = jobContext.getInMessageContext() - .getParameters(); - for (String paramKey : inputParams.keySet()) { - ActualParameter inParam = (ActualParameter) inputParams - .get(paramKey); - String paramDataType = inParam.getType().getType().toString(); - if("URI".equals(paramDataType)) { - String uri = ((URIParameterType) inParam.getType()).getValue(); - String fileName = new File(uri).getName(); + List<String> inFilePrms = extractInFileParams(); + for (String uri : inFilePrms) { + String fileName = new File(uri).getName(); if (uri.startsWith("file")) { try { String uriWithoutProtocol = uri.substring(uri.lastIndexOf("://") + 1, uri.length()); @@ -86,7 +74,6 @@ public class DataTransferrer { } } - } } } @@ -103,63 +90,16 @@ public class DataTransferrer { if(!file.exists()){ file.mkdirs(); } - - Map<String, ActualParameter> stringMap = new HashMap<String, ActualParameter>(); - - Map<String, Object> outputParams = jobContext.getOutMessageContext() - .getParameters(); - - for (String paramKey : outputParams.keySet()) { - - ActualParameter outParam = (ActualParameter) outputParams - .get(paramKey); - - String paramDataType = outParam.getType().getType().toString(); - - if ("String".equals(paramDataType)) { - String stringPrm = ((StringParameterType) outParam - .getType()).getValue(); - String localFileName = null; - //TODO: why analysis.tar? it wont scale to other gateways.. - if(stringPrm == null || stringPrm.isEmpty()){ - continue; -// localFileName = "analysis-results.tar"; - }else{ - localFileName = stringPrm.substring(stringPrm.lastIndexOf("/")+1); - } - String outputLocation = downloadLocation+File.separator+localFileName; - FileDownloader fileDownloader = new FileDownloader(stringPrm,outputLocation, Mode.overwrite); + List<String> outPrms = extractOutParams(jobContext); + for (String outPrm : outPrms) { + String outputLocation = downloadLocation+File.separator+outPrm; + FileDownloader fileDownloader = new FileDownloader(outPrm,outputLocation, Mode.overwrite); try { fileDownloader.perform(storageClient); - ((StringParameterType) outParam.getType()).setValue(outputLocation); - stringMap.put(paramKey, outParam); } catch (Exception e) { throw new GFacProviderException(e.getLocalizedMessage(),e); } - } - - else if ("StringArray".equals(paramDataType)) { - String[] valueArray = ((StringArrayType) outParam.getType()) - .getValueArray(); - for (String v : valueArray) { - String localFileName = v.substring(v.lastIndexOf("/")+1);; - String outputLocation = downloadLocation+File.separator+localFileName; - FileDownloader fileDownloader = new FileDownloader(v,outputLocation, Mode.overwrite); - try { - fileDownloader.perform(storageClient); - ((StringParameterType) outParam.getType()).setValue(outputLocation); - stringMap.put(paramKey, outParam); - } catch (Exception e) { - throw new GFacProviderException(e.getLocalizedMessage(),e); - } - } - } } - if (stringMap == null || stringMap.isEmpty()) { - log.warn("Empty Output returned from the Application, Double check the application" + - "and ApplicationDescriptor output Parameter Names"); - } - downloadStdOuts(); } @@ -171,12 +111,8 @@ public class DataTransferrer { file.mkdirs(); } - HpcApplicationDeploymentType appDepType = (HpcApplicationDeploymentType) jobContext - .getApplicationContext().getApplicationDeploymentDescription() - .getType(); - - String stdout = appDepType.getStandardOutput(); - String stderr = appDepType.getStandardError(); + String stdout = jobContext.getStandardOutput(); + String stderr = jobContext.getStandardError(); if(stdout != null) { stdout = stdout.substring(stdout.lastIndexOf('/')+1); } @@ -190,18 +126,15 @@ public class DataTransferrer { String stderrFileName = (stdout == null || stderr.equals("")) ? "stderr" : stderr; - ApplicationDescription application = jobContext.getApplicationContext().getApplicationDeploymentDescription(); - ApplicationDeploymentDescriptionType appDesc = application.getType(); - + ApplicationDeploymentDescription application = jobContext.getApplicationContext().getApplicationDeploymentDescription(); + String stdoutLocation = downloadLocation+File.separator+stdoutFileName; FileDownloader f1 = new FileDownloader(stdoutFileName,stdoutLocation, Mode.overwrite); try { f1.perform(storageClient); log.info("Downloading stdout and stderr.."); String stdoutput = readFile(stdoutLocation); - log.info("Stdout downloaded to "+stdoutLocation); - appDesc.setStandardOutput(stdoutput); - + log.info("Stdout downloaded to -> "+stdoutLocation); if(UASDataStagingProcessor.isUnicoreEndpoint(jobContext)) { String scriptExitCodeFName = "UNICORE_SCRIPT_EXIT_CODE"; String scriptCodeLocation = downloadLocation+File.separator+scriptExitCodeFName; @@ -210,52 +143,47 @@ public class DataTransferrer { f1.perform(storageClient); log.info("UNICORE_SCRIPT_EXIT_CODE downloaded to "+scriptCodeLocation); } - String stderrLocation = downloadLocation+File.separator+stderrFileName; f1.setFrom(stderrFileName); f1.setTo(stderrLocation); f1.perform(storageClient); String stderror = readFile(stderrLocation); - log.info("Stderr downloaded to "+stderrLocation); - appDesc.setStandardError(stderror); + log.info("Stderr downloaded to -> "+stderrLocation); } catch (Exception e) { throw new GFacProviderException(e.getLocalizedMessage(),e); } } - public List<String> extractOutStringParams(JobExecutionContext context) { - - Map<String, Object> outputParams = context.getOutMessageContext() - .getParameters(); - + public List<String> extractOutParams(JobExecutionContext context) { List<String> outPrmsList = new ArrayList<String>(); - - for (String paramKey : outputParams.keySet()) { - - ActualParameter outParam = (ActualParameter) outputParams - .get(paramKey); - - String paramDataType = outParam.getType().getType().toString(); - - if ("String".equals(paramDataType)) { - String strPrm = ((StringParameterType) outParam.getType()) - .getValue(); - outPrmsList.add(strPrm); - } - - else if (("StringArray").equals(paramDataType)) { - String[] uriArray = ((URIArrayType) outParam.getType()) - .getValueArray(); - for (String u : uriArray) { - outPrmsList.add(u); - } - - } - } - + List<OutputDataObjectType> applicationOutputs = jobContext.getTaskData().getApplicationOutputs(); + if (applicationOutputs != null && !applicationOutputs.isEmpty()){ + for (OutputDataObjectType output : applicationOutputs){ + if(output.getType().equals(DataType.STRING)) { + outPrmsList.add(output.getValue()); + } + else if(output.getType().equals(DataType.FLOAT) || output.getType().equals(DataType.INTEGER)) { + outPrmsList.add(String.valueOf(output.getValue())); + + } + } + } return outPrmsList; } + + public List<String> extractInFileParams() { + List<String> filePrmsList = new ArrayList<String>(); + List<InputDataObjectType> applicationInputs = jobContext.getTaskData().getApplicationInputs(); + if (applicationInputs != null && !applicationInputs.isEmpty()){ + for (InputDataObjectType output : applicationInputs){ + if(output.getType().equals(DataType.URI)) { + filePrmsList.add(output.getValue()); + } + } + } + return filePrmsList; + } private String readFile(String localFile) throws IOException { http://git-wip-us.apache.org/repos/asf/airavata/blob/28a235e8/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/JSDLGenerator.java ---------------------------------------------------------------------- diff --git a/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/JSDLGenerator.java b/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/JSDLGenerator.java index c29e12d..9755bc7 100644 --- a/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/JSDLGenerator.java +++ b/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/JSDLGenerator.java @@ -21,8 +21,8 @@ package org.apache.airavata.gfac.bes.utils; +import org.apache.airavata.gfac.core.context.ApplicationContext; import org.apache.airavata.gfac.core.context.JobExecutionContext; -import org.apache.airavata.schemas.gfac.HpcApplicationDeploymentType; import org.ggf.schemas.jsdl.x2005.x11.jsdl.JobDefinitionDocument; import org.ggf.schemas.jsdl.x2005.x11.jsdl.JobDefinitionType; import org.slf4j.Logger; @@ -41,19 +41,15 @@ public class JSDLGenerator implements BESConstants { protected final Logger log = LoggerFactory.getLogger(this.getClass()); - public synchronized static JobDefinitionDocument buildJSDLInstance( - JobExecutionContext context) throws Exception { + public synchronized static JobDefinitionDocument buildJSDLInstance(JobExecutionContext context) throws Exception { JobDefinitionDocument jobDefDoc = JobDefinitionDocument.Factory .newInstance(); JobDefinitionType value = jobDefDoc.addNewJobDefinition(); - HpcApplicationDeploymentType appDepType = (HpcApplicationDeploymentType) context - .getApplicationContext().getApplicationDeploymentDescription() - .getType(); - + // build Identification - createJobIdentification(value, appDepType); + createJobIdentification(value, context); ResourceProcessor.generateResourceElements(value, context); @@ -70,12 +66,9 @@ public class JSDLGenerator implements BESConstants { .newInstance(); JobDefinitionType value = jobDefDoc.addNewJobDefinition(); - HpcApplicationDeploymentType appDepType = (HpcApplicationDeploymentType) context - .getApplicationContext().getApplicationDeploymentDescription() - .getType(); - + // build Identification - createJobIdentification(value, appDepType); + createJobIdentification(value, context); ResourceProcessor.generateResourceElements(value, context); @@ -87,41 +80,6 @@ public class JSDLGenerator implements BESConstants { } public synchronized static JobDefinitionDocument buildJSDLInstance( - JobExecutionContext context, DataServiceInfo dataInfo) - throws Exception { - - JobDefinitionDocument jobDefDoc = JobDefinitionDocument.Factory - .newInstance(); - JobDefinitionType value = jobDefDoc.addNewJobDefinition(); - - HpcApplicationDeploymentType appDepType = (HpcApplicationDeploymentType) context - .getApplicationContext().getApplicationDeploymentDescription() - .getType(); - - createJobIdentification(value, appDepType); - - ResourceProcessor.generateResourceElements(value, context); - - ApplicationProcessor.generateJobSpecificAppElements(value, context); - - switch (dataInfo.getDirectoryAccesMode()) { - case SMSBYTEIO: - if(null == dataInfo.getDataServiceUrl() || "".equals(dataInfo.getDataServiceUrl())) - throw new Exception("No SMS address found"); - UASDataStagingProcessor.generateDataStagingElements(value, context, - dataInfo.getDataServiceUrl()); - break; - case RNSBYTEIO: - case GridFTP: - default: - DataStagingProcessor.generateDataStagingElements(value, context); - break; - - } - return jobDefDoc; - } - - public synchronized static JobDefinitionDocument buildJSDLInstance( JobExecutionContext context, String smsUrl, Object jobDirectoryMode) throws Exception { @@ -129,12 +87,8 @@ public class JSDLGenerator implements BESConstants { .newInstance(); JobDefinitionType value = jobDefDoc.addNewJobDefinition(); - HpcApplicationDeploymentType appDepType = (HpcApplicationDeploymentType) context - .getApplicationContext().getApplicationDeploymentDescription() - .getType(); - // build Identification - createJobIdentification(value, appDepType); + createJobIdentification(value, context); ResourceProcessor.generateResourceElements(value, context); @@ -146,18 +100,18 @@ public class JSDLGenerator implements BESConstants { return jobDefDoc; } - private static void createJobIdentification(JobDefinitionType value, - HpcApplicationDeploymentType appDepType) { - if (appDepType.getProjectAccount() != null) { - - if (appDepType.getProjectAccount().getProjectAccountNumber() != null) - JSDLUtils.addProjectName(value, appDepType.getProjectAccount() - .getProjectAccountNumber()); - - if (appDepType.getProjectAccount().getProjectAccountDescription() != null) - JSDLUtils.getOrCreateJobIdentification(value).setDescription( - appDepType.getProjectAccount() - .getProjectAccountDescription()); + private static void createJobIdentification(JobDefinitionType value, JobExecutionContext context) { + ApplicationContext appCtxt = context.getApplicationContext(); + + if (appCtxt != null) { + if (appCtxt.getComputeResourcePreference().getAllocationProjectNumber() != null) + JSDLUtils.addProjectName(value, appCtxt.getComputeResourcePreference().getAllocationProjectNumber()); + + if (appCtxt.getApplicationInterfaceDescription().getApplicationDescription() != null) + JSDLUtils.getOrCreateJobIdentification(value).setDescription(appCtxt.getApplicationInterfaceDescription().getApplicationDescription()); + + if (appCtxt.getApplicationInterfaceDescription().getApplicationName() != null) + JSDLUtils.getOrCreateJobIdentification(value).setJobName(appCtxt.getApplicationInterfaceDescription().getApplicationName()); } } http://git-wip-us.apache.org/repos/asf/airavata/blob/28a235e8/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/ResourceProcessor.java ---------------------------------------------------------------------- diff --git a/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/ResourceProcessor.java b/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/ResourceProcessor.java index 5df9a0f..fce0c31 100644 --- a/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/ResourceProcessor.java +++ b/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/ResourceProcessor.java @@ -25,125 +25,64 @@ import org.apache.airavata.gfac.core.context.JobExecutionContext; import org.apache.airavata.gfac.core.provider.GFacProviderException; import org.apache.airavata.model.workspace.experiment.ComputationalResourceScheduling; import org.apache.airavata.model.workspace.experiment.TaskDetails; -import org.apache.airavata.schemas.gfac.HpcApplicationDeploymentType; -import org.apache.airavata.schemas.gfac.QueueType; import org.ggf.schemas.jsdl.x2005.x11.jsdl.JobDefinitionType; -import org.ogf.schemas.jsdl.x2007.x02.jsdlSpmd.NumberOfProcessesType; public class ResourceProcessor { public static void generateResourceElements(JobDefinitionType value, JobExecutionContext context) throws Exception{ - HpcApplicationDeploymentType appDepType = (HpcApplicationDeploymentType) context - .getApplicationContext().getApplicationDeploymentDescription() - .getType(); - - createMemory(value, appDepType); TaskDetails taskData = context.getTaskData(); - if(taskData != null && taskData.isSetTaskScheduling()){ - ComputationalResourceScheduling computionResource= taskData.getTaskScheduling(); - try { - int cpuCount = computionResource.getTotalCPUCount(); - if(cpuCount>0){ -// appDepType.setCpuCount(cpuCount); - NumberOfProcessesType num = NumberOfProcessesType.Factory.newInstance(); - String processers = Integer.toString(cpuCount); - num.setStringValue(processers); - JSDLUtils.getOrCreateSPMDApplication(value).setNumberOfProcesses(num); - } - } catch (NullPointerException e) { - new GFacProviderException("No Value sent in WorkflowContextHeader for Node Count, value in the Deployment Descriptor will be used",e); - } - try { - int nodeCount = computionResource.getNodeCount(); - if(nodeCount>0){ - appDepType.setNodeCount(nodeCount); - } - } catch (NullPointerException e) { - new GFacProviderException("No Value sent in WorkflowContextHeader for Node Count, value in the Deployment Descriptor will be used",e); - } - try { - String queueName = computionResource.getQueueName(); - if (queueName != null) { - if(appDepType.getQueue() == null){ - QueueType queueType = appDepType.addNewQueue(); - queueType.setQueueName(queueName); - }else{ - appDepType.getQueue().setQueueName(queueName); - } - } - } catch (NullPointerException e) { - new GFacProviderException("No Value sent in WorkflowContextHeader for Node Count, value in the Deployment Descriptor will be used",e); - } - try { - int maxwallTime = computionResource.getWallTimeLimit(); - if(maxwallTime>0){ - appDepType.setMaxWallTime(maxwallTime); - } - } catch (NullPointerException e) { - new GFacProviderException("No Value sent in WorkflowContextHeader for Node Count, value in the Deployment Descriptor will be used",e); - } - } - if (appDepType.getCpuCount() > 0) { - RangeValueType rangeType = new RangeValueType(); - rangeType.setLowerBound(Double.NaN); - rangeType.setUpperBound(Double.NaN); - rangeType.setExact(appDepType.getCpuCount()); - JSDLUtils.setTotalCPUCountRequirements(value, rangeType); - } + if(taskData != null && taskData.isSetTaskScheduling()){ + try { + ComputationalResourceScheduling crs = taskData.getTaskScheduling(); + + if (crs.getTotalPhysicalMemory() > 0) { + RangeValueType rangeType = new RangeValueType(); + rangeType.setLowerBound(Double.NaN); + rangeType.setUpperBound(Double.NaN); + rangeType.setExact(crs.getTotalPhysicalMemory()); + JSDLUtils.setIndividualPhysicalMemoryRequirements(value, rangeType); + } + + if (crs.getNodeCount() > 0) { + RangeValueType rangeType = new RangeValueType(); + rangeType.setLowerBound(Double.NaN); + rangeType.setUpperBound(Double.NaN); + rangeType.setExact(crs.getNodeCount()); + JSDLUtils.setTotalResourceCountRequirements(value, rangeType); + } + + if(crs.getWallTimeLimit() > 0) { + RangeValueType cpuTime = new RangeValueType(); + cpuTime.setLowerBound(Double.NaN); + cpuTime.setUpperBound(Double.NaN); + long wallTime = crs.getWallTimeLimit() * 60; + cpuTime.setExact(wallTime); + JSDLUtils.setIndividualCPUTimeRequirements(value, cpuTime); + } + + if(crs.getTotalCPUCount() > 0) { + RangeValueType rangeType = new RangeValueType(); + rangeType.setLowerBound(Double.NaN); + rangeType.setUpperBound(Double.NaN); + rangeType.setExact(crs.getTotalCPUCount()); + JSDLUtils.setTotalCPUCountRequirements(value, rangeType); + } + } catch (NullPointerException npe) { + new GFacProviderException("No value set for resource requirements.",npe); + } + + + } - if (appDepType.getProcessorsPerNode() > 0) { - RangeValueType rangeType = new RangeValueType(); - rangeType.setLowerBound(Double.NaN); - rangeType.setUpperBound(Double.NaN); - rangeType.setExact(appDepType.getProcessorsPerNode()); - JSDLUtils.setIndividualCPUCountRequirements(value, rangeType); - } - if (appDepType.getNodeCount() > 0) { - RangeValueType rangeType = new RangeValueType(); - rangeType.setLowerBound(Double.NaN); - rangeType.setUpperBound(Double.NaN); - rangeType.setExact(appDepType.getNodeCount()); - JSDLUtils.setTotalResourceCountRequirements(value, rangeType); - } - - if(appDepType.getMaxWallTime() > 0) { - RangeValueType cpuTime = new RangeValueType(); - cpuTime.setLowerBound(Double.NaN); - cpuTime.setUpperBound(Double.NaN); - long wallTime = appDepType.getMaxWallTime() * 60; - cpuTime.setExact(wallTime); - JSDLUtils.setIndividualCPUTimeRequirements(value, cpuTime); - } + } - private static void createMemory(JobDefinitionType value, HpcApplicationDeploymentType appDepType){ - if (appDepType.getMinMemory() > 0 && appDepType.getMaxMemory() > 0) { - RangeValueType rangeType = new RangeValueType(); - rangeType.setLowerBound(appDepType.getMinMemory()); - rangeType.setUpperBound(appDepType.getMaxMemory()); - JSDLUtils.setIndividualPhysicalMemoryRequirements(value, rangeType); - } - - else if (appDepType.getMinMemory() > 0 && appDepType.getMaxMemory() <= 0) { - // TODO set Wall time - RangeValueType rangeType = new RangeValueType(); - rangeType.setLowerBound(appDepType.getMinMemory()); - JSDLUtils.setIndividualPhysicalMemoryRequirements(value, rangeType); - } - - else if (appDepType.getMinMemory() <= 0 && appDepType.getMaxMemory() > 0) { - // TODO set Wall time - RangeValueType rangeType = new RangeValueType(); - rangeType.setUpperBound(appDepType.getMinMemory()); - JSDLUtils.setIndividualPhysicalMemoryRequirements(value, rangeType); - } - - } + http://git-wip-us.apache.org/repos/asf/airavata/blob/28a235e8/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/UASDataStagingProcessor.java ---------------------------------------------------------------------- diff --git a/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/UASDataStagingProcessor.java b/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/UASDataStagingProcessor.java index 76624cc..9c92789 100644 --- a/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/UASDataStagingProcessor.java +++ b/modules/gfac/gfac-bes/src/main/java/org/apache/airavata/gfac/bes/utils/UASDataStagingProcessor.java @@ -22,72 +22,45 @@ package org.apache.airavata.gfac.bes.utils; import java.io.File; -import java.util.HashMap; import java.util.List; -import java.util.Map; - -import org.apache.airavata.commons.gfac.type.ActualParameter; import org.apache.airavata.gfac.core.context.JobExecutionContext; -import org.apache.airavata.gfac.core.context.MessageContext; -import org.apache.airavata.schemas.gfac.HpcApplicationDeploymentType; -import org.apache.airavata.schemas.gfac.StringArrayType; -import org.apache.airavata.schemas.gfac.StringParameterType; -import org.apache.airavata.schemas.gfac.URIArrayType; -import org.apache.airavata.schemas.gfac.URIParameterType; -import org.apache.airavata.schemas.gfac.UnicoreHostType; +import org.apache.airavata.model.appcatalog.appinterface.DataType; +import org.apache.airavata.model.appcatalog.appinterface.InputDataObjectType; +import org.apache.airavata.model.appcatalog.appinterface.OutputDataObjectType; +import org.apache.airavata.model.appcatalog.computeresource.JobSubmissionProtocol; import org.ggf.schemas.jsdl.x2005.x11.jsdl.JobDefinitionType; public class UASDataStagingProcessor { public static void generateDataStagingElements(JobDefinitionType value, JobExecutionContext context, String smsUrl) throws Exception{ - - HpcApplicationDeploymentType appDepType = (HpcApplicationDeploymentType) context - .getApplicationContext().getApplicationDeploymentDescription() - .getType(); - smsUrl = "BFT:"+smsUrl; - - if (context.getInMessageContext().getParameters().size() > 0) { - buildDataStagingFromInputContext(context, value, smsUrl, appDepType); + + if (context.getTaskData().getApplicationInputs().size() > 0) { + buildDataStagingFromInputContext(context, value, smsUrl); } -// MessageContext outMessage = new MessageContext(); -// ActualParameter a1 = new ActualParameter(); -// a1.getType().changeType(StringParameterType.type); -// ((StringParameterType)a1.getType()).setValue("analysis-results.tar"); -// outMessage.addParameter("o1", a1); -// context.setOutMessageContext(outMessage); - // now download for string typed outputs are to be done - if (context.getOutMessageContext().getParameters().size() > 0) { - buildFromOutputContext(context, value, smsUrl, appDepType); + if (context.getTaskData().getApplicationOutputs().size() > 0) { + buildFromOutputContext(context, value, smsUrl); } - //TODO need a review for us3 gateway.. -// createStdOutURIs(value, appDepType, smsUrl, isUnicoreEndpoint(context)); } - private static void createInURISMSElement(JobDefinitionType value, - String smsUrl, String inputDir, ActualParameter inParam) + private static void createInURISMSElement(JobDefinitionType value, String smsUrl, String uri) throws Exception { - - String uri = ((URIParameterType) inParam.getType()).getValue(); - //TODO: To add this input file name setting part of Airavata API String fileName = "input/" + new File(uri).getName(); if (uri.startsWith("file")) { - String fileUri = smsUrl+"#/"+fileName; - - JSDLUtils.addDataStagingSourceElement(value, fileUri, null, fileName); - } else if (uri.startsWith("gsiftp") || uri.startsWith("http") - || uri.startsWith("rns")) { - // no need to stage-in those files to the input - // directory because unicore site will fetch them for the user - JSDLUtils.addDataStagingSourceElement(value, uri, null, fileName); - } + uri = smsUrl+"#/"+fileName; + + } + // no need to stage-in those files to the input + // directory because unicore site will fetch them for the user + // supported third party transfers include + // gsiftp, http, rns, ftp + JSDLUtils.addDataStagingSourceElement(value, uri, null, fileName); } - - private static void createStdOutURIs(JobDefinitionType value, - HpcApplicationDeploymentType appDepType, String smsUrl, - boolean isUnicore) throws Exception { + + //TODO: will be deprecated + private static void createStdOutURIs(JobDefinitionType value, JobExecutionContext context, String smsUrl, boolean isUnicore) throws Exception { // no need to use smsUrl for output location, because output location is activity's working directory @@ -99,9 +72,9 @@ public class UASDataStagingProcessor { } if(!isUnicore) { - String stdout = ApplicationProcessor.getApplicationStdOut(value, appDepType); + String stdout = ApplicationProcessor.getApplicationStdOut(value, context); - String stderr = ApplicationProcessor.getApplicationStdErr(value, appDepType); + String stderr = ApplicationProcessor.getApplicationStdErr(value, context); String stdoutFileName = (stdout == null || stdout.equals("")) ? "stdout" : stdout; @@ -120,14 +93,10 @@ public class UASDataStagingProcessor { } - - private static void createOutStringElements(JobDefinitionType value, - HpcApplicationDeploymentType appDeptype, String smsUrl, String prmValue) throws Exception { - + // TODO: this should be deprecated, because the outputs are fetched using activity working dir from data transferrer + private static void createOutStringElements(JobDefinitionType value, String smsUrl, String prmValue) throws Exception { if(prmValue == null || "".equals(prmValue)) return; - String finalSMSPath = smsUrl + "#/output/"+prmValue; - JSDLUtils.addDataStagingTargetElement(value, null, prmValue, null); } @@ -140,88 +109,45 @@ public class UASDataStagingProcessor { private static JobDefinitionType buildFromOutputContext(JobExecutionContext context, - JobDefinitionType value, String smsUrl, - HpcApplicationDeploymentType appDepType) throws Exception { - - Map<String, Object> outputParams = context.getOutMessageContext() - .getParameters(); - - for (String paramKey : outputParams.keySet()) { - - ActualParameter outParam = (ActualParameter) outputParams - .get(paramKey); - - String paramDataType = outParam.getType().getType().toString(); - - if ("URI".equals(paramDataType)) { - String uriPrm = ((URIParameterType) outParam.getType()) - .getValue(); - createOutURIElement(value, uriPrm); - } - - else if (("URIArray").equals(paramDataType)) { - String[] uriArray = ((URIArrayType) outParam.getType()) - .getValueArray(); - for (String u : uriArray) { - - createOutURIElement(value, u); - } - - } -// else if ("String".equals(paramDataType)) { -// String stringPrm = ((StringParameterType) outParam -// .getType()).getValue(); -// createOutStringElements(value, appDepType, smsUrl, stringPrm); -// } -// -// else if ("StringArray".equals(paramDataType)) { -// String[] valueArray = ((StringArrayType) outParam.getType()) -// .getValueArray(); -// for (String v : valueArray) { -// createOutStringElements(value, appDepType, smsUrl, v); -// } -// } - } - + JobDefinitionType value, String smsUrl) throws Exception { + List<OutputDataObjectType> applicationOutputs = context.getTaskData().getApplicationOutputs(); + if (applicationOutputs != null && !applicationOutputs.isEmpty()){ + for (OutputDataObjectType output : applicationOutputs){ + if(output.getType().equals(DataType.URI)) { + createOutURIElement(value, output.getValue()); + } + else if(output.getType().equals(DataType.STRING)) { + // TODO: remove this check, as out string + createOutStringElements(value, smsUrl, output.getValue()); + } + } + } return value; } - private static void buildDataStagingFromInputContext(JobExecutionContext context, JobDefinitionType value, String smsUrl, HpcApplicationDeploymentType appDepType) + private static void buildDataStagingFromInputContext(JobExecutionContext context, JobDefinitionType value, String smsUrl) throws Exception { + List<InputDataObjectType> applicationInputs = context.getApplicationContext().getApplicationInterfaceDescription().getApplicationInputs(); - // TODO set data directory - Map<String, Object> inputParams = context.getInMessageContext() - .getParameters(); - - for (String paramKey : inputParams.keySet()) { - - ActualParameter inParam = (ActualParameter) inputParams - .get(paramKey); - - // if single urls then convert each url into jsdl source - // elements, - // that are formed by concat of gridftpurl+inputdir+filename - - String paramDataType = inParam.getType().getType().toString(); - - if ("URI".equals(paramDataType)) { - createInURISMSElement(value, smsUrl, - appDepType.getInputDataDirectory(), inParam); - } - - // string params are converted into the job arguments - - else if ("String".equals(paramDataType)) { - String stringPrm = ((StringParameterType) inParam.getType()) - .getValue(); - ApplicationProcessor.addApplicationArgument(value, appDepType, stringPrm); + if (applicationInputs != null && !applicationInputs.isEmpty()){ + for (InputDataObjectType input : applicationInputs){ + if(input.getType().equals(DataType.URI)){ + //TODO: set the in sms url + createInURISMSElement(value, smsUrl, input.getValue()); + } + else if(input.getType().equals(DataType.STRING)){ + ApplicationProcessor.addApplicationArgument(value, context, input.getValue()); + } + else if (input.getType().equals(DataType.FLOAT) || input.getType().equals(DataType.INTEGER)){ + ApplicationProcessor.addApplicationArgument(value, context, String.valueOf(input.getValue())); + } } } } public static boolean isUnicoreEndpoint(JobExecutionContext context) { - return ( (context.getApplicationContext().getHostDescription().getType() instanceof UnicoreHostType)?true:false ); + return context.getPreferredJobSubmissionProtocol().equals(JobSubmissionProtocol.UNICORE); } }
