diff --git a/src/prerna/engine/api/IStorageEngine.java b/src/prerna/engine/api/IStorageEngine.java index 7a7829ba92f..e3c39909642 100644 --- a/src/prerna/engine/api/IStorageEngine.java +++ b/src/prerna/engine/api/IStorageEngine.java @@ -63,6 +63,18 @@ public interface IStorageEngine extends IEngine { */ List> listDetails(String path) throws Exception; + /** + * List all versions of a specific object in storage. + * Only supported by engines with versioning enabled (S3, GCS). + * + * @param storagePath the path to the object in storage + * @return a list of version details (versionId, lastModified, size, isLatest) + * @throws Exception if listing fails + */ + default List> listVersions(String storagePath) throws Exception { + throw new UnsupportedOperationException("Object versioning is not supported by this storage engine"); + } + /** * * @param localPath @@ -81,13 +93,17 @@ public interface IStorageEngine extends IEngine { void syncStorageToLocal(String storagePath, String localPath) throws Exception; /** + * Copy local files to a storage folder path. + * Returns the version identifier if the underlying storage has versioning enabled, + * or null if versioning is not supported/enabled. * - * @param localFilePath - * @param storageFolderPath - * @param metadata - * @throws Exception + * @param localFilePath the local file or folder path(s) to upload + * @param storageFolderPath the destination path in storage + * @param metadata optional metadata to attach to the uploaded objects + * @return the version identifier (S3 versionId or GCS generation), or null if not available + * @throws Exception if the upload fails */ - void copyToStorage(String localFilePath, String storageFolderPath, Map metadata) throws Exception; + String copyToStorage(String localFilePath, String storageFolderPath, Map metadata) throws Exception; /** * @@ -141,4 +157,16 @@ default byte[] readBlobToMemory(String storagePath) throws Exception { default void updateBlobMetadata(String storagePath, Map metadata) throws Exception { throw new UnsupportedOperationException("updateBlobMetadata is not supported by this storage engine"); } + + /** + * Copy a specific version of a file from storage to local. + * + * @param storageFilePath the path to the file in storage + * @param localFolderPath the local folder to download to + * @param versionId the version identifier (S3 versionId or GCS generation number) + * @throws Exception if the operation is not supported or fails + */ + default void copyToLocal(String storageFilePath, String localFolderPath, String versionId) throws Exception { + copyToLocal(storageFilePath, localFolderPath); + } } diff --git a/src/prerna/engine/impl/storage/AWSNativeBlobStorageEngine.java b/src/prerna/engine/impl/storage/AWSNativeBlobStorageEngine.java index a6b3b37cef1..450a1c8b89e 100644 --- a/src/prerna/engine/impl/storage/AWSNativeBlobStorageEngine.java +++ b/src/prerna/engine/impl/storage/AWSNativeBlobStorageEngine.java @@ -40,11 +40,13 @@ import java.util.Collections; import java.util.Comparator; import java.util.HashMap; +import java.util.LinkedHashMap; import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Properties; import java.util.Set; +import java.util.concurrent.atomic.AtomicReference; import java.util.regex.Pattern; import java.util.stream.Collectors; import java.util.stream.Stream; @@ -68,7 +70,11 @@ import software.amazon.awssdk.services.s3.model.HeadObjectResponse; import software.amazon.awssdk.services.s3.model.ListObjectsV2Request; import software.amazon.awssdk.services.s3.model.ListObjectsV2Response; +import software.amazon.awssdk.services.s3.model.ListObjectVersionsRequest; +import software.amazon.awssdk.services.s3.model.ListObjectVersionsResponse; +import software.amazon.awssdk.services.s3.model.ObjectVersion; import software.amazon.awssdk.services.s3.model.PutObjectRequest; +import software.amazon.awssdk.services.s3.model.PutObjectResponse; import software.amazon.awssdk.services.s3.model.S3Exception; import software.amazon.awssdk.services.s3.model.S3Object; @@ -239,6 +245,39 @@ public List> listDetails(String path) throws Exception { return objectDetails; } + @Override + public List> listVersions(String storagePath) throws Exception { + List> versions = new ArrayList<>(); + String key = normalizeStoragePrefixPath(storagePath); + // Remove trailing slash for file keys + if (key.endsWith("/")) { + key = key.substring(0, key.length() - 1); + } + + ListObjectVersionsRequest request = ListObjectVersionsRequest.builder() + .bucket(this.bucket) + .prefix(key) + .build(); + + ListObjectVersionsResponse response = this.client.listObjectVersions(request); + + for (ObjectVersion version : response.versions()) { + // Only include exact key matches (not prefix matches) + if (!version.key().equals(key)) { + continue; + } + Map versionInfo = new LinkedHashMap<>(); + versionInfo.put("versionId", version.versionId()); + versionInfo.put("lastModified", version.lastModified() != null ? version.lastModified().toString() : null); + versionInfo.put("size", version.size()); + versionInfo.put("isLatest", version.isLatest()); + versionInfo.put("key", version.key()); + versions.add(versionInfo); + } + + return versions; + } + @Override public void syncLocalToStorage(String localPath, String storagePath, Map metadata) throws Exception { @@ -376,11 +415,12 @@ public void syncStorageToLocal(String storagePath, String localPath) throws Exce } @Override - public void copyToStorage(String localFilePath, String storageFolderPath, Map metadata) + public String copyToStorage(String localFilePath, String storageFolderPath, Map metadata) throws Exception { List paths = parseLocalPaths(localFilePath); List uploadedFiles = new ArrayList<>(); List failedFiles = new ArrayList<>(); + AtomicReference lastVersionId = new AtomicReference<>(null); boolean found = false; for (Path path : paths) { @@ -396,7 +436,11 @@ public void copyToStorage(String localFilePath, String storageFolderPath, Map stream = Files.walk(path)) { stream.filter(Files::isRegularFile).forEach(file -> { try { - uploadedFiles.add(uploadFile(path, file, storageFolderPath, metadata)); + String versionId = uploadFile(path, file, storageFolderPath, metadata); + uploadedFiles.add(file.toString()); + if (versionId != null) { + lastVersionId.set(versionId); + } } catch (Exception e) { failedFiles.add(file.toString()); classLogger.error("Failed to upload file: {}", file, e); @@ -407,7 +451,11 @@ public void copyToStorage(String localFilePath, String storageFolderPath, Map responseStream = this.client.getObject(getRequest)) { + Files.copy(responseStream, localFilePath, StandardCopyOption.REPLACE_EXISTING); + classLogger.info("Downloaded versioned file: {} (versionId={})", localFilePath, versionId); + } + } else { + copyToLocal(storageFilePath, localFolderPath); + } + } + @Override public void deleteFromStorage(String storagePath) throws Exception { storagePath = Utility.normalizePath(storagePath); @@ -699,9 +773,11 @@ private String uploadFileToS3(String storagePath, Path filePath, Path basePath, PutObjectRequest putRequest = PutObjectRequest.builder().bucket(this.bucket).key(fileKey).metadata(metaMap) .build(); - this.client.putObject(putRequest, filePath); + PutObjectResponse response = this.client.putObject(putRequest, filePath); classLogger.info("Uploaded/Updated file: {}", fileKey); - return; + if (response.versionId() != null) { + classLogger.info("Version ID for {}: {}", fileKey, response.versionId()); + } }, "Uploading file to S3: " + fileKey); return fileKey; @@ -871,12 +947,16 @@ private String uploadFile(Path rootPath, Path file, String storageFolderPath, Ma .collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().toString()))); } + AtomicReference versionIdRef = new AtomicReference<>(null); retryOperation(() -> { - this.client.putObject(putBuilder.build(), file); + PutObjectResponse response = this.client.putObject(putBuilder.build(), file); classLogger.info("Uploaded file to S3: {}", fileKey); + if (response.versionId() != null) { + versionIdRef.set(response.versionId()); + } }, "Uploading to S3: " + fileKey); - return fileKey; + return versionIdRef.get(); } private void downloadFile(String key, Path localFilePath) throws IOException { diff --git a/src/prerna/engine/impl/storage/AbstractRCloneStorageEngine.java b/src/prerna/engine/impl/storage/AbstractRCloneStorageEngine.java index 7f7269ed12e..531445798a9 100644 --- a/src/prerna/engine/impl/storage/AbstractRCloneStorageEngine.java +++ b/src/prerna/engine/impl/storage/AbstractRCloneStorageEngine.java @@ -135,9 +135,10 @@ public void syncStorageToLocal(String storagePath, String localPath) throws IOEx } @Override - public void copyToStorage(String localFilePath, String storageFolderPath, Map metadata) + public String copyToStorage(String localFilePath, String storageFolderPath, Map metadata) throws IOException, InterruptedException { copyToStorage(localFilePath, storageFolderPath, null, metadata); + return null; } @Override diff --git a/src/prerna/engine/impl/storage/AzureNativeBlobStorageEngine.java b/src/prerna/engine/impl/storage/AzureNativeBlobStorageEngine.java index 4e4aa895f67..0b1cffedb8e 100644 --- a/src/prerna/engine/impl/storage/AzureNativeBlobStorageEngine.java +++ b/src/prerna/engine/impl/storage/AzureNativeBlobStorageEngine.java @@ -303,7 +303,7 @@ public void syncStorageToLocal(String storagePath, String localPath) throws Exce } @Override - public void copyToStorage(String localFilePath, String storageFolderPath, Map metadata) + public String copyToStorage(String localFilePath, String storageFolderPath, Map metadata) throws Exception { // Extract container and blob directory String[] containerAndPath = extractContainerAndPath(storageFolderPath); @@ -358,6 +358,7 @@ public void copyToStorage(String localFilePath, String storageFolderPath, Map metadata) + public String copyToStorage(String localFilePath, String storageFolderPath, Map metadata) throws Exception { List sources = parseLocalPaths(localFilePath); Path destinationFolder = resolveStoragePath(storageFolderPath); @@ -302,6 +302,7 @@ public void copyToStorage(String localFilePath, String storageFolderPath, Map> listDetails(String containerPrefix) throws Exce return detailsList; } + @Override + public List> listVersions(String storagePath) throws Exception { + List> versions = new ArrayList<>(); + String key = storagePath == null ? "" : Utility.normalizePath(storagePath).trim(); + if (key.startsWith("/")) { + key = key.substring(1); + } + if (key.endsWith("/")) { + key = key.substring(0, key.length() - 1); + } + + // List all versions using versions(true) option + Page page = this.bucket.list( + Storage.BlobListOption.prefix(key), + Storage.BlobListOption.versions(true)); + + for (Blob blob : page.iterateAll()) { + // Only include exact key matches + if (!blob.getName().equals(key)) { + continue; + } + Map versionInfo = new LinkedHashMap<>(); + versionInfo.put("versionId", String.valueOf(blob.getGeneration())); + versionInfo.put("lastModified", blob.getUpdateTimeOffsetDateTime() != null + ? blob.getUpdateTimeOffsetDateTime().toString() : null); + versionInfo.put("size", blob.getSize()); + versionInfo.put("isLatest", blob.getDeleteTimeOffsetDateTime() == null); + versionInfo.put("key", blob.getName()); + versions.add(versionInfo); + } + + return versions; + } + @Override public void syncLocalToStorage(String localPath, String storagePath, Map metadata) throws Exception { @@ -375,11 +411,12 @@ public void syncStorageToLocal(String storagePath, String localPath) throws Exce } @Override - public void copyToStorage(String localFilePath, String storageFolderPath, Map metadata) + public String copyToStorage(String localFilePath, String storageFolderPath, Map metadata) throws Exception { List paths = parseLocalPaths(localFilePath); List uploadedFiles = new ArrayList<>(); List failedFiles = new ArrayList<>(); + AtomicReference lastVersionId = new AtomicReference<>(null); boolean found = false; for (Path filePath : paths) { if (!Files.exists(filePath)) { @@ -395,7 +432,11 @@ public void copyToStorage(String localFilePath, String storageFolderPath, Map stream = Files.walk(filePath)) { stream.filter(Files::isRegularFile).forEach(file -> { try { - uploadedFiles.add(uploadFileToGCS(filePath, file, storageFolderPath, metadata)); + String generation = uploadFileToGCS(filePath, file, storageFolderPath, metadata); + uploadedFiles.add(file.toString()); + if (generation != null) { + lastVersionId.set(generation); + } } catch (Exception e) { failedFiles.add(file.toString()); classLogger.error("Failed to upload file: {}", file, e); @@ -406,7 +447,11 @@ public void copyToStorage(String localFilePath, String storageFolderPath, Map deletedFiles = new ArrayList<>(); @@ -871,16 +941,20 @@ private String uploadFileToGCS(Path rootPath, Path file, String storageFolderPat .collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().toString()))); } + AtomicReference generationRef = new AtomicReference<>(null); retryOperation(() -> { try { - storage.create(blobInfoBuilder.build(), Files.readAllBytes(file)); + Blob blob = storage.create(blobInfoBuilder.build(), Files.readAllBytes(file)); classLogger.info("Uploaded file to GCS: {}", blobName); + if (blob.getGeneration() != null) { + generationRef.set(String.valueOf(blob.getGeneration())); + } } catch (IOException e) { classLogger.error("Failed to upload file to GCS: {}", blobName, e); } }, "Uploading: " + blobName); - return blobName; + return generationRef.get(); } private void deleteEmptyBlobs(String storageFolderPath) { diff --git a/src/prerna/engine/impl/storage/JCIFSStorageEngine.java b/src/prerna/engine/impl/storage/JCIFSStorageEngine.java index ab0a07ad44a..07b52864933 100644 --- a/src/prerna/engine/impl/storage/JCIFSStorageEngine.java +++ b/src/prerna/engine/impl/storage/JCIFSStorageEngine.java @@ -153,10 +153,10 @@ public void syncStorageToLocal(String storagePath, String localPath) throws Exce } @Override - public void copyToStorage(String localFilePath, String storageFolderPath, Map metadata) + public String copyToStorage(String localFilePath, String storageFolderPath, Map metadata) throws Exception { // TODO Auto-generated method stub - + return null; } @Override diff --git a/src/prerna/engine/impl/storage/SFTPStorageEngine.java b/src/prerna/engine/impl/storage/SFTPStorageEngine.java index 899d069ddc0..5334645614b 100644 --- a/src/prerna/engine/impl/storage/SFTPStorageEngine.java +++ b/src/prerna/engine/impl/storage/SFTPStorageEngine.java @@ -272,7 +272,7 @@ public void syncStorageToLocal(String storagePath, String localPath) throws Exce } @Override - public void copyToStorage(String localFilePath, String storageFolderPath, Map metadata) + public String copyToStorage(String localFilePath, String storageFolderPath, Map metadata) throws Exception { SSHClient sshClient = null; SFTPClient sftpClient = null; @@ -306,6 +306,7 @@ public void copyToStorage(String localFilePath, String storageFolderPath, MapPixel usage: + *
+ * ListStorageVersions(storage=["engineId"], storagePath=["path/to/file.pdf"]);
+ * 
+ * + *

Returns a list of version details including versionId, lastModified, size, and isLatest. + * Only supported by storage engines with versioning enabled (AWS S3, GCS). + */ +public class ListStorageVersionsReactor extends AbstractReactor { + + private static final Logger classLogger = LogManager.getLogger(ListStorageVersionsReactor.class); + + public ListStorageVersionsReactor() { + this.keysToGet = new String[] { ReactorKeysEnum.STORAGE.getKey(), ReactorKeysEnum.STORAGE_PATH.getKey() }; + this.keyRequired = new int[] { 1, 1 }; + } + + @Override + public NounMetadata execute() { + organizeKeys(); + IStorageEngine storage = getStorage(); + String storagePath = this.keyValue.get(ReactorKeysEnum.STORAGE_PATH.getKey()); + + try { + List> versions = storage.listVersions(storagePath); + return new NounMetadata(versions, PixelDataType.VECTOR); + } catch (UnsupportedOperationException e) { + throw new IllegalArgumentException("This storage engine does not support version listing"); + } catch (Exception e) { + classLogger.error("Error listing versions for path: {}", storagePath, e); + throw new IllegalArgumentException("Error listing storage versions at path: " + storagePath); + } + } + + private IStorageEngine getStorage() { + GenRowStruct grs = this.store.getGenRowStruct(ReactorKeysEnum.STORAGE.getKey()); + if (grs != null && !grs.isEmpty()) { + IStorageEngine storage = null; + if (grs.get(0) instanceof String) { + String storageId = (String) grs.get(0); + if (!SecurityEngineUtils.userCanViewEngine(this.insight.getUser(), storageId)) { + throw new IllegalArgumentException( + "Storage " + storageId + " does not exist or user does not have access to storage"); + } + storage = Utility.getStorage(storageId); + } else { + storage = (IStorageEngine) grs.get(0); + } + return storage; + } + + List storageInputs = this.curRow.getNounsOfType(PixelDataType.STORAGE); + if (storageInputs != null && !storageInputs.isEmpty()) { + return (IStorageEngine) storageInputs.get(0).getValue(); + } + + throw new NullPointerException("No storage engine defined"); + } +} diff --git a/src/prerna/reactor/storage/PullFromStorageReactor.java b/src/prerna/reactor/storage/PullFromStorageReactor.java index bf3d2ba67bd..13bb6e0e75e 100644 --- a/src/prerna/reactor/storage/PullFromStorageReactor.java +++ b/src/prerna/reactor/storage/PullFromStorageReactor.java @@ -44,14 +44,29 @@ import prerna.util.UploadInputUtility; import prerna.util.Utility; +/** + * Pull files from a storage path to a local path. + * + * Pixel usage: PullFromStorage(storage=[""], storagePath=[""], filePath=[""], version=[""]); + * + * Parameters: + * storage (String, required) - The storage engine instance or id + * storagePath (String, required) - The storage path(s) to download from + * filePath (String, required) - The local path to download files to + * space (String, optional) - The project space context + * version (String, optional) - The version ID to download a specific object version + * + * Returns: BOOLEAN - true on success. + */ public class PullFromStorageReactor extends AbstractReactor { private static final Logger classLogger = LogManager.getLogger(PullFromStorageReactor.class); public PullFromStorageReactor() { this.keysToGet = new String[] { ReactorKeysEnum.STORAGE.getKey(), ReactorKeysEnum.STORAGE_PATH.getKey(), - ReactorKeysEnum.SPACE.getKey(), ReactorKeysEnum.FILE_PATH.getKey() }; - this.keyRequired = new int[] { 1, 1, 0, 1 }; + ReactorKeysEnum.SPACE.getKey(), ReactorKeysEnum.FILE_PATH.getKey(), + ReactorKeysEnum.VERSION.getKey() }; + this.keyRequired = new int[] { 1, 1, 0, 1, 0 }; } @Override @@ -64,8 +79,14 @@ public NounMetadata execute() { new File(fileLocation).mkdirs(); } + String versionId = this.keyValue.get(ReactorKeysEnum.VERSION.getKey()); + try { - storage.copyToLocal(storagePath, fileLocation); + if (versionId != null && !versionId.isEmpty()) { + storage.copyToLocal(storagePath, fileLocation, versionId); + } else { + storage.copyToLocal(storagePath, fileLocation); + } return new NounMetadata(true, PixelDataType.BOOLEAN); } catch (Exception e) { classLogger.error(Constants.STACKTRACE, e); @@ -111,6 +132,8 @@ protected String getDescriptionForKey(String key) { return "The storage path(s) to download from"; } else if (key.equals(ReactorKeysEnum.FILE_PATH.getKey())) { return "The local path to download files to"; + } else if (key.equals(ReactorKeysEnum.VERSION.getKey())) { + return "Optional version ID to download a specific object version (S3 versionId or GCS generation)"; } return super.getDescriptionForKey(key); } diff --git a/src/prerna/reactor/storage/PushToStorageReactor.java b/src/prerna/reactor/storage/PushToStorageReactor.java index e7a86bd42a0..cda2dd30fbc 100644 --- a/src/prerna/reactor/storage/PushToStorageReactor.java +++ b/src/prerna/reactor/storage/PushToStorageReactor.java @@ -28,6 +28,7 @@ package prerna.reactor.storage; import java.io.File; +import java.util.HashMap; import java.util.List; import java.util.Map; @@ -45,6 +46,20 @@ import prerna.util.UploadInputUtility; import prerna.util.Utility; +/** + * Push files from a local path to a storage path. + * + * Pixel usage: PushToStorage(storage=[""], storagePath=[""], filePath=[""], metadata=[{}]); + * + * Parameters: + * storage (String, required) - The storage engine instance or id + * storagePath (String, required) - The storage path to upload files to + * filePath (String, required) - The local path(s) to upload from + * space (String, optional) - The project space context + * metadata (Map, optional) - Metadata to attach to uploaded objects + * + * Returns: MAP - containing success status and versionId (if versioning enabled), or BOOLEAN true. + */ public class PushToStorageReactor extends AbstractReactor { private static final Logger classLogger = LogManager.getLogger(PushToStorageReactor.class); @@ -71,7 +86,13 @@ public NounMetadata execute() { Map metadata = getMetadata(); try { - storage.copyToStorage(fileLocation, storageFolderPath, metadata); + String versionId = storage.copyToStorage(fileLocation, storageFolderPath, metadata); + if (versionId != null && !versionId.isEmpty()) { + Map result = new HashMap<>(); + result.put("success", true); + result.put("versionId", versionId); + return new NounMetadata(result, PixelDataType.MAP); + } return new NounMetadata(true, PixelDataType.BOOLEAN); } catch (Exception e) { classLogger.error(Constants.STACKTRACE, e); @@ -128,6 +149,8 @@ protected String getDescriptionForKey(String key) { return "The storage path to upload files to"; } else if (key.equals(ReactorKeysEnum.FILE_PATH.getKey())) { return "The local path(s) to upload from"; + } else if (key.equals(ReactorKeysEnum.METADATA.getKey())) { + return "Optional metadata map to attach to uploaded objects"; } return super.getDescriptionForKey(key); }