This is an automated email from the ASF dual-hosted git repository.
dimuthuupe pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/airavata-mft.git
The following commit(s) were added to refs/heads/develop by this push:
new 7c866be Fixed https://github.com/apache/airavata-mft/issues/38
7c866be is described below
commit 7c866be0a23ad49ae629d93916177a050080b9b4
Author: Dimuthu Wannipurage <[email protected]>
AuthorDate: Wed May 19 01:55:17 2021 -0400
Fixed https://github.com/apache/airavata-mft/issues/38
---
.../airavata/mft/agent/TransportMediator.java | 4 +-
.../airavata/mft/agent/http/HttpServerHandler.java | 9 ++-
.../mft/agent/http/HttpTransferRequest.java | 10 +++
.../apache/airavata/mft/agent/rpc/RPCParser.java | 2 +
api/service/pom.xml | 6 ++
.../airavata/mft/api/handler/MFTApiHandler.java | 1 +
api/stub/src/main/proto/MFTApi.proto | 40 ++++++------
.../org/apache/airavata/mft/core/TransferTask.java | 10 ++-
.../apache/airavata/mft/core/api/Connector.java | 1 +
pom.xml | 2 +-
.../mft/transport/azure/AzureReceiver.java | 6 ++
.../airavata/mft/transport/azure/AzureSender.java | 6 ++
.../airavata/mft/transport/box/BoxReceiver.java | 6 ++
.../airavata/mft/transport/box/BoxSender.java | 6 ++
.../mft/transport/dropbox/DropboxReceiver.java | 6 ++
.../mft/transport/dropbox/DropboxSender.java | 6 ++
.../airavata/mft/transport/ftp/FTPReceiver.java | 6 ++
.../airavata/mft/transport/ftp/FTPSender.java | 6 ++
.../airavata/mft/transport/gcp/GCSReceiver.java | 9 ++-
.../airavata/mft/transport/gcp/GCSSender.java | 11 +++-
.../mft/transport/local/LocalReceiver.java | 6 ++
.../airavata/mft/transport/local/LocalSender.java | 7 +++
.../airavata/mft/transport/s3/S3Receiver.java | 6 ++
.../apache/airavata/mft/transport/s3/S3Sender.java | 6 ++
.../mft/transport/scp/SCPMetadataCollector.java | 71 ++++++++++++++++++++--
.../airavata/mft/transport/scp/SCPReceiver.java | 37 ++++++++++-
.../airavata/mft/transport/scp/SCPSender.java | 39 ++++++++++--
27 files changed, 282 insertions(+), 43 deletions(-)
diff --git
a/agent/src/main/java/org/apache/airavata/mft/agent/TransportMediator.java
b/agent/src/main/java/org/apache/airavata/mft/agent/TransportMediator.java
index 40bdd0d..b9195b8 100644
--- a/agent/src/main/java/org/apache/airavata/mft/agent/TransportMediator.java
+++ b/agent/src/main/java/org/apache/airavata/mft/agent/TransportMediator.java
@@ -69,9 +69,9 @@ public class TransportMediator {
context.setTransferId(transferId);
TransferTask recvTask = new
TransferTask(request.getMftAuthorizationToken(), request.getSourceResourceId(),
- request.getSourceToken(), context, inConnector);
+ request.getSourceChildResourcePath(),
request.getSourceToken(), context, inConnector);
TransferTask sendTask = new
TransferTask(request.getMftAuthorizationToken(),
request.getDestinationResourceId(),
- request.getDestinationToken(), context, outConnector);
+ request.getDestinationChildResourcePath(),
request.getDestinationToken(), context, outConnector);
List<Future<Integer>> futureList = new ArrayList<>();
ExecutorCompletionService<Integer> completionService = new
ExecutorCompletionService<>(executor);
diff --git
a/agent/src/main/java/org/apache/airavata/mft/agent/http/HttpServerHandler.java
b/agent/src/main/java/org/apache/airavata/mft/agent/http/HttpServerHandler.java
index cda8af5..1941d54 100644
---
a/agent/src/main/java/org/apache/airavata/mft/agent/http/HttpServerHandler.java
+++
b/agent/src/main/java/org/apache/airavata/mft/agent/http/HttpServerHandler.java
@@ -98,6 +98,7 @@ public class HttpServerHandler extends
SimpleChannelInboundHandler<FullHttpReque
FileResourceMetadata fileResourceMetadata =
metadataCollector.getFileResourceMetadata(authToken,
httpTransferRequest.getResourceId(),
+ httpTransferRequest.getChildResourcePath(),
httpTransferRequest.getCredentialToken());
long fileLength = fileResourceMetadata.getResourceSize();
@@ -122,7 +123,9 @@ public class HttpServerHandler extends
SimpleChannelInboundHandler<FullHttpReque
connectorContext.setTransferId(uri);
connectorContext.setMetadata(new FileResourceMetadata()); // TODO
Resolve
- TransferTask pullTask = new TransferTask(authToken,
httpTransferRequest.getResourceId(), httpTransferRequest.getCredentialToken(),
connectorContext, connector);
+ TransferTask pullTask = new TransferTask(authToken,
httpTransferRequest.getResourceId(),
+ httpTransferRequest.getChildResourcePath(),
httpTransferRequest.getCredentialToken(),
+ connectorContext, connector);
// TODO aggregate pullStatusFuture and sendFileFuture for keepalive
test
Future<Integer> pullStatusFuture = executor.submit(pullTask);
@@ -137,9 +140,9 @@ public class HttpServerHandler extends
SimpleChannelInboundHandler<FullHttpReque
@Override
public void operationProgressed(ChannelProgressiveFuture future,
long progress, long total) {
if (total < 0) { // total unknown
- System.err.println(future.channel() + " Transfer progress:
" + progress);
+ logger.error(future.channel() + " Transfer progress: " +
progress);
} else {
- System.err.println(future.channel() + " Transfer progress:
" + progress + " / " + total);
+ logger.error(future.channel() + " Transfer progress: " +
progress + " / " + total);
}
}
diff --git
a/agent/src/main/java/org/apache/airavata/mft/agent/http/HttpTransferRequest.java
b/agent/src/main/java/org/apache/airavata/mft/agent/http/HttpTransferRequest.java
index f030acd..1633ab4 100644
---
a/agent/src/main/java/org/apache/airavata/mft/agent/http/HttpTransferRequest.java
+++
b/agent/src/main/java/org/apache/airavata/mft/agent/http/HttpTransferRequest.java
@@ -25,6 +25,7 @@ public class HttpTransferRequest {
private MetadataCollector otherMetadataCollector;
private ConnectorParams connectorParams;
private String resourceId;
+ private String childResourcePath;
private String credentialToken;
private long createdTime = System.currentTimeMillis();
@@ -55,6 +56,15 @@ public class HttpTransferRequest {
return this;
}
+ public String getChildResourcePath() {
+ return childResourcePath;
+ }
+
+ public HttpTransferRequest setChildResourcePath(String childResourcePath) {
+ this.childResourcePath = childResourcePath;
+ return this;
+ }
+
public String getCredentialToken() {
return credentialToken;
}
diff --git
a/agent/src/main/java/org/apache/airavata/mft/agent/rpc/RPCParser.java
b/agent/src/main/java/org/apache/airavata/mft/agent/rpc/RPCParser.java
index f4be490..a58a89b 100644
--- a/agent/src/main/java/org/apache/airavata/mft/agent/rpc/RPCParser.java
+++ b/agent/src/main/java/org/apache/airavata/mft/agent/rpc/RPCParser.java
@@ -150,6 +150,7 @@ public class RPCParser {
case "submitHttpDownload":
resourceId = request.getParameters().get("resourceId");
+ String childResourcePath =
request.getParameters().get("childResourcePath");
String sourceToken =
request.getParameters().get("sourceToken");
String storeType = request.getParameters().get("storeType");
@@ -168,6 +169,7 @@ public class RPCParser {
.setSecretServiceHost(secretServiceHost)
.setSecretServicePort(secretServicePort));
transferRequest.setResourceId(resourceId);
+ transferRequest.setChildResourcePath(childResourcePath);
transferRequest.setCredentialToken(sourceToken);
transferRequest.setOtherMetadataCollector(metadataCollectorOp.get());
transferRequest.setOtherConnector(connectorOp.get());
diff --git a/api/service/pom.xml b/api/service/pom.xml
index ae93a01..7f4459f 100644
--- a/api/service/pom.xml
+++ b/api/service/pom.xml
@@ -70,6 +70,12 @@
<groupId>com.google.protobuf</groupId>
<artifactId>protobuf-java-util</artifactId>
<version>3.15.0</version>
+ <exclusions>
+ <exclusion>
+ <groupId>com.google.guava</groupId>
+ <artifactId>guava</artifactId>
+ </exclusion>
+ </exclusions>
</dependency>
<!-- To be removed -->
diff --git
a/api/service/src/main/java/org/apache/airavata/mft/api/handler/MFTApiHandler.java
b/api/service/src/main/java/org/apache/airavata/mft/api/handler/MFTApiHandler.java
index c6005d2..c3a9840 100644
---
a/api/service/src/main/java/org/apache/airavata/mft/api/handler/MFTApiHandler.java
+++
b/api/service/src/main/java/org/apache/airavata/mft/api/handler/MFTApiHandler.java
@@ -102,6 +102,7 @@ public class MFTApiHandler extends
MFTApiServiceGrpc.MFTApiServiceImplBase {
.withMessageId(UUID.randomUUID().toString())
.withMethod("submitHttpDownload")
.withParameter("resourceId", request.getSourceResourceId())
+ .withParameter("childResourcePath",
request.getSourceResourceChildPath())
.withParameter("sourceToken", request.getSourceToken())
.withParameter("storeType", request.getSourceType())
.withParameter("mftAuthorizationToken",
JsonFormat.printer().print(request.getMftAuthorizationToken()));
diff --git a/api/stub/src/main/proto/MFTApi.proto
b/api/stub/src/main/proto/MFTApi.proto
index 0370619..66db86f 100644
--- a/api/stub/src/main/proto/MFTApi.proto
+++ b/api/stub/src/main/proto/MFTApi.proto
@@ -9,14 +9,16 @@ import "CredCommon.proto";
message TransferApiRequest {
string sourceResourceId = 1;
- string sourceType = 2;
- string sourceToken = 3;
- string destinationResourceId = 7;
- string destinationType = 9;
- string destinationToken = 10;
- bool affinityTransfer = 13;
- map<string, int32> targetAgents = 14;
- org.apache.airavata.mft.common.AuthToken mftAuthorizationToken = 15;
+ string sourceChildResourcePath = 2;
+ string sourceType = 3;
+ string sourceToken = 4;
+ string destinationResourceId = 5;
+ string destinationChildResourcePath = 6;
+ string destinationType = 7;
+ string destinationToken = 8;
+ bool affinityTransfer = 9;
+ map<string, int32> targetAgents = 10;
+ org.apache.airavata.mft.common.AuthToken mftAuthorizationToken = 11;
}
message TransferApiResponse {
@@ -25,6 +27,7 @@ message TransferApiResponse {
message HttpUploadApiRequest {
string destinationResourceId = 1;
+ string destinationResourceChildPath = 2;
string destinationToken = 3;
string destinationType = 4;
string targetAgent = 5;
@@ -38,6 +41,7 @@ message HttpUploadApiResponse {
message HttpDownloadApiRequest {
string sourceResourceId = 1;
+ string sourceResourceChildPath = 2;
string sourceToken = 3;
string sourceType = 4;
string targetAgent = 5;
@@ -63,11 +67,12 @@ message TransferStateApiResponse {
message ResourceAvailabilityRequest {
string resourceId = 1;
- string resourceType = 2;
- string resourceToken = 3;
- string resourceBackend = 4;
- string resourceCredentialBackend = 5;
- org.apache.airavata.mft.common.AuthToken mftAuthorizationToken = 6;
+ string childResourcePath = 2;
+ string resourceType = 3;
+ string resourceToken = 4;
+ string resourceBackend = 5;
+ string resourceCredentialBackend = 6;
+ org.apache.airavata.mft.common.AuthToken mftAuthorizationToken = 7;
}
message ResourceAvailabilityResponse {
@@ -99,10 +104,11 @@ message DirectoryMetadataResponse {
message FetchResourceMetadataRequest {
string resourceId = 1;
- string resourceType = 2;
- string resourceToken = 3;
- string resourceBackend = 4;
- string resourceCredentialBackend = 5;
+ string childResourcePath = 2;
+ string resourceType = 3;
+ string resourceToken = 4;
+ string resourceBackend = 5;
+ string resourceCredentialBackend = 6;
string targetAgentId = 7;
string childPath = 8; // if the child entities of the parent resource are
required, set this field
org.apache.airavata.mft.common.AuthToken mftAuthorizationToken = 9;
diff --git a/core/src/main/java/org/apache/airavata/mft/core/TransferTask.java
b/core/src/main/java/org/apache/airavata/mft/core/TransferTask.java
index 38440f6..42a0e08 100644
--- a/core/src/main/java/org/apache/airavata/mft/core/TransferTask.java
+++ b/core/src/main/java/org/apache/airavata/mft/core/TransferTask.java
@@ -27,21 +27,27 @@ public class TransferTask implements Callable<Integer> {
private Connector connector;
private ConnectorContext context;
private String resourceId;
+ private String childResourcePath;
private String credentialToken;
private AuthToken authToken;
- public TransferTask(AuthToken authToken, String resourceId, String
credentialToken,
+ public TransferTask(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken,
ConnectorContext context, Connector connector) {
this.connector = connector;
this.context = context;
this.resourceId = resourceId;
this.authToken = authToken;
this.credentialToken = credentialToken;
+ this.childResourcePath = childResourcePath;
}
@Override
public Integer call() throws Exception {
- this.connector.startStream(authToken, resourceId, credentialToken,
context);
+ if (childResourcePath == null || "".equals(childResourcePath)) {
+ this.connector.startStream(authToken, resourceId, credentialToken,
context);
+ } else {
+ this.connector.startStream(authToken, resourceId,
childResourcePath, credentialToken, context);
+ }
return 0;
}
}
diff --git a/core/src/main/java/org/apache/airavata/mft/core/api/Connector.java
b/core/src/main/java/org/apache/airavata/mft/core/api/Connector.java
index 88df5b1..b7b5e61 100644
--- a/core/src/main/java/org/apache/airavata/mft/core/api/Connector.java
+++ b/core/src/main/java/org/apache/airavata/mft/core/api/Connector.java
@@ -25,4 +25,5 @@ public interface Connector {
String secretServiceHost, int secretServicePort) throws
Exception;
public void destroy();
void startStream(AuthToken authToken, String resourceId, String
credentialToken, ConnectorContext context) throws Exception;
+ void startStream(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken, ConnectorContext context) throws
Exception;
}
diff --git a/pom.xml b/pom.xml
index 7c3921c..57ce5a3 100755
--- a/pom.xml
+++ b/pom.xml
@@ -43,12 +43,12 @@
<modules>
<module>common</module>
<module>core</module>
+ <module>api</module>
<module>transport</module>
<module>agent</module>
<module>services</module>
<module>admin</module>
<module>controller</module>
- <module>api</module>
<module>examples</module>
</modules>
diff --git
a/transport/azure-transport/src/main/java/org/apache/airavata/mft/transport/azure/AzureReceiver.java
b/transport/azure-transport/src/main/java/org/apache/airavata/mft/transport/azure/AzureReceiver.java
index 3c09756..229b319 100644
---
a/transport/azure-transport/src/main/java/org/apache/airavata/mft/transport/azure/AzureReceiver.java
+++
b/transport/azure-transport/src/main/java/org/apache/airavata/mft/transport/azure/AzureReceiver.java
@@ -126,4 +126,10 @@ public class AzureReceiver implements Connector {
streamOs.close();
logger.info("Completed azure receive for remote server for transfer
{}", context.getTransferId());
}
+
+ @Override
+ public void startStream(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken,
+ ConnectorContext context) throws Exception {
+ throw new UnsupportedOperationException();
+ }
}
diff --git
a/transport/azure-transport/src/main/java/org/apache/airavata/mft/transport/azure/AzureSender.java
b/transport/azure-transport/src/main/java/org/apache/airavata/mft/transport/azure/AzureSender.java
index 53dd84a..2e7457a 100644
---
a/transport/azure-transport/src/main/java/org/apache/airavata/mft/transport/azure/AzureSender.java
+++
b/transport/azure-transport/src/main/java/org/apache/airavata/mft/transport/azure/AzureSender.java
@@ -94,4 +94,10 @@ public class AzureSender implements Connector {
blockBlobClient.upload(context.getStreamBuffer().getInputStream(),
context.getMetadata().getResourceSize(), true);
logger.info("Completed Azure send for remote server for transfer {}",
context.getTransferId());
}
+
+ @Override
+ public void startStream(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken,
+ ConnectorContext context) throws Exception {
+ throw new UnsupportedOperationException();
+ }
}
diff --git
a/transport/box-transport/src/main/java/org/apache/airavata/mft/transport/box/BoxReceiver.java
b/transport/box-transport/src/main/java/org/apache/airavata/mft/transport/box/BoxReceiver.java
index a68fbcf..94a735e 100644
---
a/transport/box-transport/src/main/java/org/apache/airavata/mft/transport/box/BoxReceiver.java
+++
b/transport/box-transport/src/main/java/org/apache/airavata/mft/transport/box/BoxReceiver.java
@@ -91,4 +91,10 @@ public class BoxReceiver implements Connector {
logger.info("Completed Box Receiver stream for transfer {}",
context.getTransferId());
}
+
+ @Override
+ public void startStream(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken,
+ ConnectorContext context) throws Exception {
+ throw new UnsupportedOperationException();
+ }
}
diff --git
a/transport/box-transport/src/main/java/org/apache/airavata/mft/transport/box/BoxSender.java
b/transport/box-transport/src/main/java/org/apache/airavata/mft/transport/box/BoxSender.java
index 700a99f..946234a 100644
---
a/transport/box-transport/src/main/java/org/apache/airavata/mft/transport/box/BoxSender.java
+++
b/transport/box-transport/src/main/java/org/apache/airavata/mft/transport/box/BoxSender.java
@@ -91,4 +91,10 @@ public class BoxSender implements Connector {
logger.info("Completed Box Sender stream for transfer {}",
context.getTransferId());
}
+
+ @Override
+ public void startStream(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken,
+ ConnectorContext context) throws Exception {
+ throw new UnsupportedOperationException();
+ }
}
diff --git
a/transport/dropbox-transport/src/main/java/org/apache/airavata/mft/transport/dropbox/DropboxReceiver.java
b/transport/dropbox-transport/src/main/java/org/apache/airavata/mft/transport/dropbox/DropboxReceiver.java
index 4215fa5..d696913 100644
---
a/transport/dropbox-transport/src/main/java/org/apache/airavata/mft/transport/dropbox/DropboxReceiver.java
+++
b/transport/dropbox-transport/src/main/java/org/apache/airavata/mft/transport/dropbox/DropboxReceiver.java
@@ -114,4 +114,10 @@ public class DropboxReceiver implements Connector {
logger.info("Completed Dropbox Receiver stream for transfer {}",
context.getTransferId());
}
+
+ @Override
+ public void startStream(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken,
+ ConnectorContext context) throws Exception {
+ throw new UnsupportedOperationException();
+ }
}
diff --git
a/transport/dropbox-transport/src/main/java/org/apache/airavata/mft/transport/dropbox/DropboxSender.java
b/transport/dropbox-transport/src/main/java/org/apache/airavata/mft/transport/dropbox/DropboxSender.java
index 0eeae00..f2ed3fc 100644
---
a/transport/dropbox-transport/src/main/java/org/apache/airavata/mft/transport/dropbox/DropboxSender.java
+++
b/transport/dropbox-transport/src/main/java/org/apache/airavata/mft/transport/dropbox/DropboxSender.java
@@ -88,4 +88,10 @@ public class DropboxSender implements Connector {
}
+
+ @Override
+ public void startStream(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken,
+ ConnectorContext context) throws Exception {
+ throw new UnsupportedOperationException();
+ }
}
diff --git
a/transport/ftp-transport/src/main/java/org/apache/airavata/mft/transport/ftp/FTPReceiver.java
b/transport/ftp-transport/src/main/java/org/apache/airavata/mft/transport/ftp/FTPReceiver.java
index f68803f..53b2620 100644
---
a/transport/ftp-transport/src/main/java/org/apache/airavata/mft/transport/ftp/FTPReceiver.java
+++
b/transport/ftp-transport/src/main/java/org/apache/airavata/mft/transport/ftp/FTPReceiver.java
@@ -126,4 +126,10 @@ public class FTPReceiver implements Connector {
throw new IllegalStateException("FTP Receiver is not initialized");
}
}
+
+ @Override
+ public void startStream(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken,
+ ConnectorContext context) throws Exception {
+ throw new UnsupportedOperationException();
+ }
}
diff --git
a/transport/ftp-transport/src/main/java/org/apache/airavata/mft/transport/ftp/FTPSender.java
b/transport/ftp-transport/src/main/java/org/apache/airavata/mft/transport/ftp/FTPSender.java
index 1ff62a8..e62bab9 100644
---
a/transport/ftp-transport/src/main/java/org/apache/airavata/mft/transport/ftp/FTPSender.java
+++
b/transport/ftp-transport/src/main/java/org/apache/airavata/mft/transport/ftp/FTPSender.java
@@ -125,4 +125,10 @@ public class FTPSender implements Connector {
throw new IllegalStateException("FTP Sender is not initialized");
}
}
+
+ @Override
+ public void startStream(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken,
+ ConnectorContext context) throws Exception {
+ throw new UnsupportedOperationException();
+ }
}
diff --git
a/transport/gcp-transport/src/main/java/org/apache/airavata/mft/transport/gcp/GCSReceiver.java
b/transport/gcp-transport/src/main/java/org/apache/airavata/mft/transport/gcp/GCSReceiver.java
index 065425f..15ade66 100644
---
a/transport/gcp-transport/src/main/java/org/apache/airavata/mft/transport/gcp/GCSReceiver.java
+++
b/transport/gcp-transport/src/main/java/org/apache/airavata/mft/transport/gcp/GCSReceiver.java
@@ -31,12 +31,9 @@ import
org.apache.airavata.mft.credential.stubs.gcs.GCSSecret;
import org.apache.airavata.mft.credential.stubs.gcs.GCSSecretGetRequest;
import org.apache.airavata.mft.resource.client.ResourceServiceClient;
import org.apache.airavata.mft.resource.client.ResourceServiceClientBuilder;
-import org.apache.airavata.mft.resource.client.StorageServiceClient;
-import org.apache.airavata.mft.resource.client.StorageServiceClientBuilder;
import org.apache.airavata.mft.resource.stubs.common.GenericResource;
import org.apache.airavata.mft.resource.stubs.common.GenericResourceGetRequest;
import org.apache.airavata.mft.resource.stubs.gcs.storage.GCSStorage;
-import org.apache.airavata.mft.resource.stubs.gcs.storage.GCSStorageGetRequest;
import org.apache.airavata.mft.secret.client.SecretServiceClient;
import org.apache.airavata.mft.secret.client.SecretServiceClientBuilder;
import org.slf4j.Logger;
@@ -137,4 +134,10 @@ public class GCSReceiver implements Connector {
logger.info("Completed GCS Receiver stream for transfer {}",
context.getTransferId());
}
+
+ @Override
+ public void startStream(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken,
+ ConnectorContext context) throws Exception {
+ throw new UnsupportedOperationException();
+ }
}
diff --git
a/transport/gcp-transport/src/main/java/org/apache/airavata/mft/transport/gcp/GCSSender.java
b/transport/gcp-transport/src/main/java/org/apache/airavata/mft/transport/gcp/GCSSender.java
index 9ae2193..7e91fb2 100644
---
a/transport/gcp-transport/src/main/java/org/apache/airavata/mft/transport/gcp/GCSSender.java
+++
b/transport/gcp-transport/src/main/java/org/apache/airavata/mft/transport/gcp/GCSSender.java
@@ -37,12 +37,9 @@ import
org.apache.airavata.mft.credential.stubs.gcs.GCSSecret;
import org.apache.airavata.mft.credential.stubs.gcs.GCSSecretGetRequest;
import org.apache.airavata.mft.resource.client.ResourceServiceClient;
import org.apache.airavata.mft.resource.client.ResourceServiceClientBuilder;
-import org.apache.airavata.mft.resource.client.StorageServiceClient;
-import org.apache.airavata.mft.resource.client.StorageServiceClientBuilder;
import org.apache.airavata.mft.resource.stubs.common.GenericResource;
import org.apache.airavata.mft.resource.stubs.common.GenericResourceGetRequest;
import org.apache.airavata.mft.resource.stubs.gcs.storage.GCSStorage;
-import org.apache.airavata.mft.resource.stubs.gcs.storage.GCSStorageGetRequest;
import org.apache.airavata.mft.secret.client.SecretServiceClient;
import org.apache.airavata.mft.secret.client.SecretServiceClientBuilder;
import org.slf4j.Logger;
@@ -127,4 +124,12 @@ public class GCSSender implements Connector {
logger.info("Completed GCS Sender stream for transfer {}",
context.getTransferId());
}
+
+ @Override
+ public void startStream(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken,
+ ConnectorContext context) throws Exception {
+ throw new UnsupportedOperationException();
+ }
+
+
}
diff --git
a/transport/local-transport/src/main/java/org/apache/airavata/mft/transport/local/LocalReceiver.java
b/transport/local-transport/src/main/java/org/apache/airavata/mft/transport/local/LocalReceiver.java
index 89718b4..84d7746 100644
---
a/transport/local-transport/src/main/java/org/apache/airavata/mft/transport/local/LocalReceiver.java
+++
b/transport/local-transport/src/main/java/org/apache/airavata/mft/transport/local/LocalReceiver.java
@@ -110,4 +110,10 @@ public class LocalReceiver implements Connector {
streamOs.close();
logger.info("Completed local receiver stream for transfer {}",
context.getTransferId());
}
+
+ @Override
+ public void startStream(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken,
+ ConnectorContext context) throws Exception {
+ throw new UnsupportedOperationException();
+ }
}
diff --git
a/transport/local-transport/src/main/java/org/apache/airavata/mft/transport/local/LocalSender.java
b/transport/local-transport/src/main/java/org/apache/airavata/mft/transport/local/LocalSender.java
index 85ef4ef..7b7e659 100644
---
a/transport/local-transport/src/main/java/org/apache/airavata/mft/transport/local/LocalSender.java
+++
b/transport/local-transport/src/main/java/org/apache/airavata/mft/transport/local/LocalSender.java
@@ -111,4 +111,11 @@ public class LocalSender implements Connector {
logger.info("Completed local sender stream for transfer {}",
context.getTransferId());
}
+
+ @Override
+ public void startStream(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken,
+ ConnectorContext context) throws Exception {
+ throw new UnsupportedOperationException();
+ }
+
}
diff --git
a/transport/s3-transport/src/main/java/org/apache/airavata/mft/transport/s3/S3Receiver.java
b/transport/s3-transport/src/main/java/org/apache/airavata/mft/transport/s3/S3Receiver.java
index 094f684..a430323 100644
---
a/transport/s3-transport/src/main/java/org/apache/airavata/mft/transport/s3/S3Receiver.java
+++
b/transport/s3-transport/src/main/java/org/apache/airavata/mft/transport/s3/S3Receiver.java
@@ -106,4 +106,10 @@ public class S3Receiver implements Connector {
logger.info("Completed S3 Receiver stream for transfer {}",
context.getTransferId());
}
+
+ @Override
+ public void startStream(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken,
+ ConnectorContext context) throws Exception {
+ throw new UnsupportedOperationException();
+ }
}
diff --git
a/transport/s3-transport/src/main/java/org/apache/airavata/mft/transport/s3/S3Sender.java
b/transport/s3-transport/src/main/java/org/apache/airavata/mft/transport/s3/S3Sender.java
index a749e21..0dec96d 100644
---
a/transport/s3-transport/src/main/java/org/apache/airavata/mft/transport/s3/S3Sender.java
+++
b/transport/s3-transport/src/main/java/org/apache/airavata/mft/transport/s3/S3Sender.java
@@ -95,4 +95,10 @@ public class S3Sender implements Connector {
context.getStreamBuffer().getInputStream(), metadata);
logger.info("Completed S3 Sender stream for transfer {}",
context.getTransferId());
}
+
+ @Override
+ public void startStream(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken,
+ ConnectorContext context) throws Exception {
+ throw new UnsupportedOperationException();
+ }
}
diff --git
a/transport/scp-transport/src/main/java/org/apache/airavata/mft/transport/scp/SCPMetadataCollector.java
b/transport/scp-transport/src/main/java/org/apache/airavata/mft/transport/scp/SCPMetadataCollector.java
index 35ad63b..0465777 100644
---
a/transport/scp-transport/src/main/java/org/apache/airavata/mft/transport/scp/SCPMetadataCollector.java
+++
b/transport/scp-transport/src/main/java/org/apache/airavata/mft/transport/scp/SCPMetadataCollector.java
@@ -131,17 +131,47 @@ public class SCPMetadataCollector implements
MetadataCollector {
}
@Override
- public FileResourceMetadata getFileResourceMetadata(AuthToken authZToken,
String parentResourceId, String resourcePath, String credentialToken) throws
Exception {
+ public FileResourceMetadata getFileResourceMetadata(AuthToken authZToken,
String parentResourceId, String childResourcePath, String credentialToken)
throws Exception {
ResourceServiceClient resourceClient =
ResourceServiceClientBuilder.buildClient(resourceServiceHost,
resourceServicePort);
- GenericResource scpResource =
resourceClient.get().getGenericResource(GenericResourceGetRequest.newBuilder().setResourceId(parentResourceId).build());
+ GenericResource resource =
resourceClient.get().getGenericResource(GenericResourceGetRequest.newBuilder().setResourceId(parentResourceId).build());
SecretServiceClient secretClient =
SecretServiceClientBuilder.buildClient(secretServiceHost, secretServicePort);
SCPSecret scpSecret =
secretClient.scp().getSCPSecret(SCPSecretGetRequest.newBuilder().setSecretId(credentialToken).build());
+ String resourcePath = null;
+
+ switch (resource.getResourceCase()){
+ case FILE:
+ resourcePath = resource.getFile().getResourcePath();
+ break;
+ case DIRECTORY:
+ resourcePath = resource.getDirectory().getResourcePath();
+ break;
+ case RESOURCE_NOT_SET:
+ throw new Exception("Resource was not set in resource with id
" + parentResourceId);
+ }
+
+ if (childResourcePath != null && !"".equals(childResourcePath)) {
+ if (resourcePath.startsWith("/")) {
+ // Linux
+ resourcePath = resourcePath.endsWith("/") ?
+ resourcePath + childResourcePath : resourcePath + "/"
+ childResourcePath;
+ } else if (resourcePath.contains("\\")) {
+ // Windows
+ resourcePath = resourcePath.endsWith("\\") ?
+ resourcePath + childResourcePath : resourcePath + "\\"
+ childResourcePath;
+ } else {
+ logger.error("Couldn't detect path seperator to append child
path {} resource path {}",
+ childResourcePath, resourcePath);
+ throw new Exception("Couldn't detect path seperator to append
child path " + childResourcePath
+ +" resource path " + resourcePath);
+ }
+ }
+
GenericResource scpResource2 = GenericResource.newBuilder()
.setFile(FileResource.newBuilder()
.setResourcePath(resourcePath).build())
-
.setScpStorage(scpResource.getScpStorage()).build();
+
.setScpStorage(resource.getScpStorage()).build();
return getFileResourceMetadata(authZToken, scpResource2, scpSecret);
}
@@ -203,16 +233,45 @@ public class SCPMetadataCollector implements
MetadataCollector {
}
@Override
- public DirectoryResourceMetadata getDirectoryResourceMetadata(AuthToken
authZToken, String parentResourceId, String resourcePath, String
credentialToken) throws Exception {
+ public DirectoryResourceMetadata getDirectoryResourceMetadata(AuthToken
authZToken, String parentResourceId, String childResourcePath, String
credentialToken) throws Exception {
ResourceServiceClient resourceClient =
ResourceServiceClientBuilder.buildClient(resourceServiceHost,
resourceServicePort);
- GenericResource scpPResource =
resourceClient.get().getGenericResource(GenericResourceGetRequest.newBuilder().setResourceId(parentResourceId).build());
+ GenericResource resource =
resourceClient.get().getGenericResource(GenericResourceGetRequest.newBuilder().setResourceId(parentResourceId).build());
SecretServiceClient secretClient =
SecretServiceClientBuilder.buildClient(secretServiceHost, secretServicePort);
SCPSecret scpSecret =
secretClient.scp().getSCPSecret(SCPSecretGetRequest.newBuilder().setSecretId(credentialToken).build());
+ String resourcePath = null;
+
+ switch (resource.getResourceCase()){
+ case FILE:
+ resourcePath = resource.getFile().getResourcePath();
+ break;
+ case DIRECTORY:
+ resourcePath = resource.getDirectory().getResourcePath();
+ case RESOURCE_NOT_SET:
+ throw new Exception("Resource was not set in resource with id
" + parentResourceId);
+ }
+
+ if (childResourcePath != null && !"".equals(childResourcePath)) {
+ if (resourcePath.startsWith("/")) {
+ // Linux
+ resourcePath = resourcePath.endsWith("/") ?
+ resourcePath + childResourcePath : resourcePath + "/"
+ childResourcePath;
+ } else if (resourcePath.contains("\\")) {
+ // Windows
+ resourcePath = resourcePath.endsWith("\\") ?
+ resourcePath + childResourcePath : resourcePath + "\\"
+ childResourcePath;
+ } else {
+ logger.error("Couldn't detect path seperator to append child
path {} resource path {}",
+ childResourcePath, resourcePath);
+ throw new Exception("Couldn't detect path seperator to append
child path " + childResourcePath
+ +" resource path " + resourcePath);
+ }
+ }
+
GenericResource scpResource = GenericResource.newBuilder()
.setDirectory(DirectoryResource.newBuilder().setResourcePath(resourcePath).build())
- .setScpStorage(scpPResource.getScpStorage()).build();
+ .setScpStorage(resource.getScpStorage()).build();
return getDirectoryResourceMetadata(authZToken,scpResource, scpSecret);
}
diff --git
a/transport/scp-transport/src/main/java/org/apache/airavata/mft/transport/scp/SCPReceiver.java
b/transport/scp-transport/src/main/java/org/apache/airavata/mft/transport/scp/SCPReceiver.java
index a3cacdf..233f3b7 100644
---
a/transport/scp-transport/src/main/java/org/apache/airavata/mft/transport/scp/SCPReceiver.java
+++
b/transport/scp-transport/src/main/java/org/apache/airavata/mft/transport/scp/SCPReceiver.java
@@ -81,6 +81,12 @@ public class SCPReceiver implements Connector {
}
public void startStream(AuthToken authToken, String resourceId, String
credentialToken, ConnectorContext context) throws Exception {
+ startStream(authToken, resourceId, null, credentialToken, context);
+ }
+
+ @Override
+ public void startStream(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken,
+ ConnectorContext context) throws Exception {
checkInitialized();
ResourceServiceClient resourceClient =
ResourceServiceClientBuilder.buildClient(resourceServiceHost,
resourceServicePort);
@@ -113,9 +119,36 @@ public class SCPReceiver implements Connector {
throw new Exception("Session can not be null. Make sure that SCP
Receiver is properly initialized");
}
- transferRemoteToStream(session, resource.getFile().getResourcePath(),
context.getStreamBuffer());
- logger.info("SCP Receive completed. Transfer {}",
context.getTransferId());
+ String resourcePath = null;
+ switch (resource.getResourceCase()){
+ case FILE:
+ resourcePath = resource.getFile().getResourcePath();
+ break;
+ case DIRECTORY:
+ resourcePath = resource.getDirectory().getResourcePath();
+ break;
+ case RESOURCE_NOT_SET:
+ throw new Exception("Resource was not set in resource with id
" + resourceId);
+ }
+ if (childResourcePath != null && !"".equals(childResourcePath)) {
+ if (resourcePath.startsWith("/")) {
+ // Linux
+ resourcePath = resourcePath.endsWith("/") ?
+ resourcePath + childResourcePath : resourcePath + "/"
+ childResourcePath;
+ } else if (resourcePath.contains("\\")) {
+ // Windows
+ resourcePath = resourcePath.endsWith("\\") ?
+ resourcePath + childResourcePath : resourcePath + "\\"
+ childResourcePath;
+ } else {
+ logger.error("Couldn't detect path seperator to append child
path {} resource path {}",
+ childResourcePath, resourcePath);
+ throw new Exception("Couldn't detect path seperator to append
child path " + childResourcePath
+ +" resource path " + resourcePath);
+ }
+ }
+ transferRemoteToStream(session, resourcePath,
context.getStreamBuffer());
+ logger.info("SCP Receive completed. Transfer {}",
context.getTransferId());
}
private void transferRemoteToStream(Session session, String from,
DoubleStreamingBuffer streamBuffer) throws Exception {
diff --git
a/transport/scp-transport/src/main/java/org/apache/airavata/mft/transport/scp/SCPSender.java
b/transport/scp-transport/src/main/java/org/apache/airavata/mft/transport/scp/SCPSender.java
index 429610f..4df432c 100644
---
a/transport/scp-transport/src/main/java/org/apache/airavata/mft/transport/scp/SCPSender.java
+++
b/transport/scp-transport/src/main/java/org/apache/airavata/mft/transport/scp/SCPSender.java
@@ -28,12 +28,9 @@ import
org.apache.airavata.mft.credential.stubs.scp.SCPSecret;
import org.apache.airavata.mft.credential.stubs.scp.SCPSecretGetRequest;
import org.apache.airavata.mft.resource.client.ResourceServiceClient;
import org.apache.airavata.mft.resource.client.ResourceServiceClientBuilder;
-import org.apache.airavata.mft.resource.client.StorageServiceClient;
-import org.apache.airavata.mft.resource.client.StorageServiceClientBuilder;
import org.apache.airavata.mft.resource.stubs.common.GenericResource;
import org.apache.airavata.mft.resource.stubs.common.GenericResourceGetRequest;
import org.apache.airavata.mft.resource.stubs.scp.storage.SCPStorage;
-import org.apache.airavata.mft.resource.stubs.scp.storage.SCPStorageGetRequest;
import org.apache.airavata.mft.secret.client.SecretServiceClient;
import org.apache.airavata.mft.secret.client.SecretServiceClientBuilder;
import org.slf4j.Logger;
@@ -84,7 +81,11 @@ public class SCPSender implements Connector {
}
public void startStream(AuthToken authToken, String resourceId, String
credentialToken, ConnectorContext context) throws Exception {
+ startStream(authToken, resourceId, null, credentialToken, context);
+ }
+ @Override
+ public void startStream(AuthToken authToken, String resourceId, String
childResourcePath, String credentialToken, ConnectorContext context) throws
Exception {
checkInitialized();
ResourceServiceClient resourceClient =
ResourceServiceClientBuilder.buildClient(resourceServiceHost,
resourceServicePort);
@@ -116,9 +117,39 @@ public class SCPSender implements Connector {
System.out.println("Session can not be null. Make sure that SCP
Sender is properly initialized");
throw new Exception("Session can not be null. Make sure that SCP
Sender is properly initialized");
}
+
+ String resourcePath = null;
+ switch (resource.getResourceCase()){
+ case FILE:
+ resourcePath = resource.getFile().getResourcePath();
+ break;
+ case DIRECTORY:
+ resourcePath = resource.getDirectory().getResourcePath();
+ break;
+ case RESOURCE_NOT_SET:
+ throw new Exception("Resource was not set in resource with id
" + resourceId);
+ }
+
+ if (childResourcePath != null && !"".equals(childResourcePath)) {
+ if (resourcePath.startsWith("/")) {
+ // Linux
+ resourcePath = resourcePath.endsWith("/") ?
+ resourcePath + childResourcePath : resourcePath + "/"
+ childResourcePath;
+ } else if (resourcePath.contains("\\")) {
+ // Windows
+ resourcePath = resourcePath.endsWith("\\") ?
+ resourcePath + childResourcePath : resourcePath + "\\"
+ childResourcePath;
+ } else {
+ logger.error("Couldn't detect path seperator to append child
path {} resource path {}",
+ childResourcePath, resourcePath);
+ throw new Exception("Couldn't detect path seperator to append
child path " + childResourcePath
+ +" resource path " + resourcePath);
+ }
+ }
+
try {
copyLocalToRemote(this.session,
- resource.getFile().getResourcePath(),
+ resourcePath,
context.getStreamBuffer(),
context.getMetadata().getResourceSize());
logger.info("SCP send to transfer {} completed",
context.getTransferId());