From 51e55c140553891650dd586c681731116cc30bb7 Mon Sep 17 00:00:00 2001 From: "Patel, Parth" Date: Thu, 11 Jun 2026 10:05:28 -0400 Subject: [PATCH 1/6] feat: allow for object versioining in gcp aws --- src/prerna/engine/api/IStorageEngine.java | 38 +++++- .../storage/AWSNativeBlobStorageEngine.java | 128 +++++++++++++++++- .../impl/storage/AbstractStorageEngine.java | 17 +++ .../GoogleCloudNativeBlobStorageEngine.java | 125 +++++++++++++++++ .../storage/PullFromStorageReactor.java | 29 +++- .../reactor/storage/PushToStorageReactor.java | 32 ++++- 6 files changed, 358 insertions(+), 11 deletions(-) diff --git a/src/prerna/engine/api/IStorageEngine.java b/src/prerna/engine/api/IStorageEngine.java index 7a7829ba92f..38bc57abe29 100644 --- a/src/prerna/engine/api/IStorageEngine.java +++ b/src/prerna/engine/api/IStorageEngine.java @@ -81,14 +81,32 @@ public interface IStorageEngine extends IEngine { void syncStorageToLocal(String storagePath, String localPath) throws Exception; /** + * Copy local files to a storage folder path. * - * @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 + * @throws Exception if the upload fails */ void copyToStorage(String localFilePath, String storageFolderPath, Map metadata) throws Exception; + /** + * Copy local files to storage and return the version identifier of the uploaded object. + * Only supported by engines with versioning enabled (S3, GCS). + * Default implementation delegates to {@link #copyToStorage} and returns null. + * + * @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 supported + * @throws Exception if the upload fails + */ + default String copyToStorageVersioned(String localFilePath, String storageFolderPath, + Map metadata) throws Exception { + copyToStorage(localFilePath, storageFolderPath, metadata); + return null; + } + /** * * @param storageFilePath @@ -141,4 +159,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..27200185ae3 100644 --- a/src/prerna/engine/impl/storage/AWSNativeBlobStorageEngine.java +++ b/src/prerna/engine/impl/storage/AWSNativeBlobStorageEngine.java @@ -45,6 +45,7 @@ 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; @@ -69,6 +70,7 @@ import software.amazon.awssdk.services.s3.model.ListObjectsV2Request; import software.amazon.awssdk.services.s3.model.ListObjectsV2Response; 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; @@ -416,7 +418,7 @@ public void 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) { + if (!Files.exists(path)) { + classLogger.error("File not found: {}", path); + failedFiles.add(path.toString()); + continue; + } + + deleteEmptyDirectories(path); + + if (Files.isDirectory(path)) { + try (Stream stream = Files.walk(path)) { + stream.filter(Files::isRegularFile).forEach(file -> { + try { + String versionId = uploadFileVersioned(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); + rollbackUploads(this.client, failedFiles); + } + }); + found = true; + } + } else { + try { + String versionId = uploadFileVersioned(path.getParent(), path, storageFolderPath, metadata); + uploadedFiles.add(path.toString()); + if (versionId != null) { + lastVersionId.set(versionId); + } + found = true; + } catch (Exception e) { + failedFiles.add(path.toString()); + classLogger.error("Failed to upload file: {}", path, e); + rollbackUploads(this.client, failedFiles); + } + } + } + // Delete empty blobs from S3 + deleteEmptyBlobsFromS3(storageFolderPath); + if (uploadedFiles.isEmpty()) { + classLogger.info("No files were uploaded."); + } else { + classLogger.info("Successfully uploaded files: {}", uploadedFiles); + } + classLogger.info(found ? "Copy completed successfully for: {}" : "No files found to copy for: {}", + storageFolderPath); + return lastVersionId.get(); + } + @Override public void copyToLocal(String storageFilePath, String localFolderPath) throws Exception { List paths = parseStorageObjectPaths(storageFilePath); @@ -510,6 +574,31 @@ public void copyToLocal(String storageFilePath, String localFolderPath) throws E storageFilePath); } + @Override + public void copyToLocal(String storageFilePath, String localFolderPath, String versionId) throws Exception { + if (versionId != null && !versionId.isEmpty()) { + String key = normalizeStoragePrefixPath(storageFilePath); + Path localDirectory = Paths.get(localFolderPath); + Files.createDirectories(localDirectory); + + String fileName = key.contains("/") ? key.substring(key.lastIndexOf("/") + 1) : key; + Path localFilePath = localDirectory.resolve(fileName); + + GetObjectRequest getRequest = GetObjectRequest.builder() + .bucket(this.bucket) + .key(key) + .versionId(versionId) + .build(); + + try (ResponseInputStream 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 +788,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 (this.versioningEnabled && response.versionId() != null) { + classLogger.info("Version ID for {}: {}", fileKey, response.versionId()); + } }, "Uploading file to S3: " + fileKey); return fileKey; @@ -879,6 +970,37 @@ private String uploadFile(Path rootPath, Path file, String storageFolderPath, Ma return fileKey; } + private String uploadFileVersioned(Path rootPath, Path file, String storageFolderPath, + Map metadata) throws IOException { + String normalizedPath = Utility.normalizePath(storageFolderPath).trim(); + if (normalizedPath.startsWith("/")) { + normalizedPath = normalizedPath.substring(1); + } + + String relativePath = Utility.normalizePath(rootPath.relativize(file).toString()).trim(); + String fileKey = normalizedPath.isEmpty() ? relativePath + : (normalizedPath.endsWith("/") ? normalizedPath + relativePath : normalizedPath + "/" + relativePath); + + PutObjectRequest.Builder putBuilder = PutObjectRequest.builder().bucket(this.bucket).key(fileKey); + + if (metadata != null && !metadata.isEmpty()) { + putBuilder.metadata(metadata.entrySet().stream() + .collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().toString()))); + } + + AtomicReference versionIdRef = new AtomicReference<>(null); + retryOperation(() -> { + PutObjectResponse response = this.client.putObject(putBuilder.build(), file); + classLogger.info("Uploaded file to S3: {}", fileKey); + if (response.versionId() != null) { + versionIdRef.set(response.versionId()); + classLogger.info("Version ID for {}: {}", fileKey, response.versionId()); + } + }, "Uploading to S3: " + fileKey); + + return versionIdRef.get(); + } + private void downloadFile(String key, Path localFilePath) throws IOException { GetObjectRequest getObjectRequest = GetObjectRequest.builder().bucket(this.bucket).key(key).build(); diff --git a/src/prerna/engine/impl/storage/AbstractStorageEngine.java b/src/prerna/engine/impl/storage/AbstractStorageEngine.java index 979c77eadc7..51f4226a307 100644 --- a/src/prerna/engine/impl/storage/AbstractStorageEngine.java +++ b/src/prerna/engine/impl/storage/AbstractStorageEngine.java @@ -47,6 +47,10 @@ public abstract class AbstractStorageEngine extends AbstractEngine implements IS protected static final Gson GSON = new GsonBuilder().disableHtmlEscaping() .setObjectToNumberStrategy(ToNumberPolicy.LONG_OR_DOUBLE).create(); + public static final String VERSIONING_ENABLED_KEY = "VERSIONING_ENABLED"; + + protected boolean versioningEnabled = false; + /** * Init the general storage values * @@ -56,6 +60,19 @@ public abstract class AbstractStorageEngine extends AbstractEngine implements IS @Override public void open(Properties smssProp) throws Exception { super.open(smssProp); + String versioningStr = smssProp.getProperty(VERSIONING_ENABLED_KEY); + if (versioningStr != null && !versioningStr.isEmpty()) { + this.versioningEnabled = Boolean.parseBoolean(versioningStr); + } + } + + /** + * Check if versioning is enabled for this storage engine. + * + * @return true if versioning is enabled + */ + public boolean isVersioningEnabled() { + return this.versioningEnabled; } /** diff --git a/src/prerna/engine/impl/storage/GoogleCloudNativeBlobStorageEngine.java b/src/prerna/engine/impl/storage/GoogleCloudNativeBlobStorageEngine.java index ea418708935..d90aee6e74e 100644 --- a/src/prerna/engine/impl/storage/GoogleCloudNativeBlobStorageEngine.java +++ b/src/prerna/engine/impl/storage/GoogleCloudNativeBlobStorageEngine.java @@ -48,6 +48,7 @@ import java.util.Map; import java.util.Properties; import java.util.Set; +import java.util.concurrent.atomic.AtomicReference; import java.util.stream.Collectors; import java.util.stream.Stream; @@ -428,6 +429,71 @@ public void 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)) { + classLogger.error("File not found: {}", filePath); + failedFiles.add(filePath.toString()); + continue; + } + + // Delete empty directories before upload + deleteEmptyDirectories(filePath); + + if (Files.isDirectory(filePath)) { + try (Stream stream = Files.walk(filePath)) { + stream.filter(Files::isRegularFile).forEach(file -> { + try { + String generation = uploadFileToGCSVersioned(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); + rollbackUploads(storage, failedFiles); + } + }); + found = true; + } + } else { + try { + String generation = uploadFileToGCSVersioned(filePath.getParent(), filePath, storageFolderPath, + metadata); + uploadedFiles.add(filePath.toString()); + if (generation != null) { + lastVersionId.set(generation); + } + found = true; + } catch (Exception e) { + failedFiles.add(filePath.toString()); + classLogger.error("Failed to upload file: {}", filePath, e); + rollbackUploads(storage, failedFiles); + } + } + } + + // Delete empty blobs from GCS + deleteEmptyBlobs(storageFolderPath); + + if (uploadedFiles.isEmpty()) { + classLogger.info("No files were uploaded."); + } else { + classLogger.info("Successfully uploaded files: {}", uploadedFiles); + } + classLogger.info(found ? "Copy completed successfully for: {}" : "No files found to copy for: {}", + storageFolderPath); + return lastVersionId.get(); + } + @Override public void copyToLocal(String storageFilePath, String localFolderPath) throws Exception { List paths = parseStorageObjectPaths(storageFilePath); @@ -487,6 +553,30 @@ public void copyToLocal(String storageFilePath, String localFolderPath) throws E storageFilePath); } + @Override + public void copyToLocal(String storageFilePath, String localFolderPath, String versionId) throws Exception { + if (versionId != null && !versionId.isEmpty()) { + String key = normalizeStoragePrefixPath(storageFilePath); + Path localDirectory = Paths.get(localFolderPath); + Files.createDirectories(localDirectory); + + String fileName = key.contains("/") ? key.substring(key.lastIndexOf("/") + 1) : key; + Path localFilePath = localDirectory.resolve(fileName); + + BlobId blobId = BlobId.of(this.BUCKET, key, Long.parseLong(versionId)); + Blob blob = storage.get(blobId); + if (blob == null) { + throw new IllegalArgumentException( + "Object not found in GCS: " + key + " with generation=" + versionId); + } + + downloadFile(blob, localFilePath); + classLogger.info("Downloaded versioned file: {} (generation={})", localFilePath, versionId); + } else { + copyToLocal(storageFilePath, localFolderPath); + } + } + @Override public void deleteFromStorage(String storagePath) throws Exception { List deletedFiles = new ArrayList<>(); @@ -883,6 +973,41 @@ private String uploadFileToGCS(Path rootPath, Path file, String storageFolderPat return blobName; } + private String uploadFileToGCSVersioned(Path rootPath, Path file, String storageFolderPath, + Map metadata) throws IOException { + String normalizedPath = Utility.normalizePath(storageFolderPath).trim(); + if (normalizedPath.startsWith("/")) { + normalizedPath = normalizedPath.substring(1); + } + String relativePath = Utility.normalizePath(rootPath.relativize(file).toString()).trim(); + String blobName = normalizedPath.isEmpty() ? relativePath + : (normalizedPath.endsWith("/") ? normalizedPath + relativePath : normalizedPath + "/" + relativePath); + + BlobId blobId = BlobId.of(this.BUCKET, blobName); + BlobInfo.Builder blobInfoBuilder = BlobInfo.newBuilder(blobId); + + if (metadata != null && !metadata.isEmpty()) { + blobInfoBuilder.setMetadata(metadata.entrySet().stream() + .collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().toString()))); + } + + AtomicReference generationRef = new AtomicReference<>(null); + retryOperation(() -> { + try { + 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())); + classLogger.info("Generation for {}: {}", blobName, blob.getGeneration()); + } + } catch (IOException e) { + classLogger.error("Failed to upload file to GCS: {}", blobName, e); + } + }, "Uploading: " + blobName); + + return generationRef.get(); + } + private void deleteEmptyBlobs(String storageFolderPath) { Page blobs = storage.list(this.BUCKET, Storage.BlobListOption.prefix(storageFolderPath)); 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..a6bf262ce75 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; @@ -36,6 +37,7 @@ import prerna.auth.utils.SecurityEngineUtils; import prerna.engine.api.IStorageEngine; +import prerna.engine.impl.storage.AbstractStorageEngine; import prerna.reactor.AbstractReactor; import prerna.sablecc2.om.GenRowStruct; import prerna.sablecc2.om.PixelDataType; @@ -45,6 +47,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 +87,19 @@ public NounMetadata execute() { Map metadata = getMetadata(); try { - storage.copyToStorage(fileLocation, storageFolderPath, metadata); + // If this engine supports versioning, use the versioned upload + if (storage instanceof AbstractStorageEngine + && ((AbstractStorageEngine) storage).isVersioningEnabled()) { + String versionId = storage.copyToStorageVersioned(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); + } + } else { + storage.copyToStorage(fileLocation, storageFolderPath, metadata); + } return new NounMetadata(true, PixelDataType.BOOLEAN); } catch (Exception e) { classLogger.error(Constants.STACKTRACE, e); @@ -128,6 +156,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); } From 73e941db3da99511d18b05e10da0656e6785e960 Mon Sep 17 00:00:00 2001 From: "Patel, Parth" Date: Thu, 11 Jun 2026 11:21:09 -0400 Subject: [PATCH 2/6] fix: move to Interface --- src/prerna/engine/api/IStorageEngine.java | 9 +++++++++ .../engine/impl/storage/AbstractStorageEngine.java | 1 + src/prerna/reactor/storage/PushToStorageReactor.java | 4 +--- 3 files changed, 11 insertions(+), 3 deletions(-) diff --git a/src/prerna/engine/api/IStorageEngine.java b/src/prerna/engine/api/IStorageEngine.java index 38bc57abe29..149a5d15456 100644 --- a/src/prerna/engine/api/IStorageEngine.java +++ b/src/prerna/engine/api/IStorageEngine.java @@ -90,6 +90,15 @@ public interface IStorageEngine extends IEngine { */ void copyToStorage(String localFilePath, String storageFolderPath, Map metadata) throws Exception; + /** + * Check if object versioning is enabled for this storage engine. + * + * @return true if versioning is enabled, false by default + */ + default boolean isVersioningEnabled() { + return false; + } + /** * Copy local files to storage and return the version identifier of the uploaded object. * Only supported by engines with versioning enabled (S3, GCS). diff --git a/src/prerna/engine/impl/storage/AbstractStorageEngine.java b/src/prerna/engine/impl/storage/AbstractStorageEngine.java index 51f4226a307..9fb04d04c4c 100644 --- a/src/prerna/engine/impl/storage/AbstractStorageEngine.java +++ b/src/prerna/engine/impl/storage/AbstractStorageEngine.java @@ -71,6 +71,7 @@ public void open(Properties smssProp) throws Exception { * * @return true if versioning is enabled */ + @Override public boolean isVersioningEnabled() { return this.versioningEnabled; } diff --git a/src/prerna/reactor/storage/PushToStorageReactor.java b/src/prerna/reactor/storage/PushToStorageReactor.java index a6bf262ce75..05aeb0d65ba 100644 --- a/src/prerna/reactor/storage/PushToStorageReactor.java +++ b/src/prerna/reactor/storage/PushToStorageReactor.java @@ -37,7 +37,6 @@ import prerna.auth.utils.SecurityEngineUtils; import prerna.engine.api.IStorageEngine; -import prerna.engine.impl.storage.AbstractStorageEngine; import prerna.reactor.AbstractReactor; import prerna.sablecc2.om.GenRowStruct; import prerna.sablecc2.om.PixelDataType; @@ -88,8 +87,7 @@ public NounMetadata execute() { Map metadata = getMetadata(); try { // If this engine supports versioning, use the versioned upload - if (storage instanceof AbstractStorageEngine - && ((AbstractStorageEngine) storage).isVersioningEnabled()) { + if (storage.isVersioningEnabled()) { String versionId = storage.copyToStorageVersioned(fileLocation, storageFolderPath, metadata); if (versionId != null && !versionId.isEmpty()) { Map result = new HashMap<>(); From 3cf65421a628e0de28cf193c4987334ecfdb8445 Mon Sep 17 00:00:00 2001 From: "Patel, Parth" Date: Thu, 11 Jun 2026 11:36:00 -0400 Subject: [PATCH 3/6] feat: add new listStorageVersionsReactor for gcp/aws --- src/prerna/engine/api/IStorageEngine.java | 12 ++++++ .../storage/AWSNativeBlobStorageEngine.java | 37 +++++++++++++++++++ .../GoogleCloudNativeBlobStorageEngine.java | 35 ++++++++++++++++++ 3 files changed, 84 insertions(+) diff --git a/src/prerna/engine/api/IStorageEngine.java b/src/prerna/engine/api/IStorageEngine.java index 149a5d15456..62952e74d59 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 diff --git a/src/prerna/engine/impl/storage/AWSNativeBlobStorageEngine.java b/src/prerna/engine/impl/storage/AWSNativeBlobStorageEngine.java index 27200185ae3..cdbe227e94b 100644 --- a/src/prerna/engine/impl/storage/AWSNativeBlobStorageEngine.java +++ b/src/prerna/engine/impl/storage/AWSNativeBlobStorageEngine.java @@ -40,6 +40,7 @@ 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; @@ -69,6 +70,9 @@ 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; @@ -241,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 { diff --git a/src/prerna/engine/impl/storage/GoogleCloudNativeBlobStorageEngine.java b/src/prerna/engine/impl/storage/GoogleCloudNativeBlobStorageEngine.java index d90aee6e74e..ddbb8f2f2d0 100644 --- a/src/prerna/engine/impl/storage/GoogleCloudNativeBlobStorageEngine.java +++ b/src/prerna/engine/impl/storage/GoogleCloudNativeBlobStorageEngine.java @@ -44,6 +44,7 @@ import java.util.Comparator; import java.util.HashMap; import java.util.HashSet; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Properties; @@ -249,6 +250,40 @@ public List> 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.isDirectory()); + versionInfo.put("key", blob.getName()); + versions.add(versionInfo); + } + + return versions; + } + @Override public void syncLocalToStorage(String localPath, String storagePath, Map metadata) throws Exception { From 5cc2de137589e67396859bbcafa2717e9ceae20b Mon Sep 17 00:00:00 2001 From: "Patel, Parth" Date: Thu, 11 Jun 2026 11:43:50 -0400 Subject: [PATCH 4/6] chore: actually commit file --- .../storage/ListStorageVersionsReactor.java | 110 ++++++++++++++++++ 1 file changed, 110 insertions(+) create mode 100644 src/prerna/reactor/storage/ListStorageVersionsReactor.java diff --git a/src/prerna/reactor/storage/ListStorageVersionsReactor.java b/src/prerna/reactor/storage/ListStorageVersionsReactor.java new file mode 100644 index 00000000000..1b45b8c8278 --- /dev/null +++ b/src/prerna/reactor/storage/ListStorageVersionsReactor.java @@ -0,0 +1,110 @@ +/******************************************************************************* + * Copyright 2015 Defense Health Agency (DHA) + * + * If your use of this software does not include any GPLv2 components: + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * ---------------------------------------------------------------------------- + * If your use of this software includes any GPLv2 components: + * This program is free software; you can redistribute it and/or + * modify it under the terms of the GNU General Public License + * as published by the Free Software Foundation; either version 2 + * of the License, or (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU General Public License for more details. + *******************************************************************************/ +package prerna.reactor.storage; + +import java.util.List; +import java.util.Map; + +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import prerna.auth.utils.SecurityEngineUtils; +import prerna.engine.api.IStorageEngine; +import prerna.reactor.AbstractReactor; +import prerna.sablecc2.om.GenRowStruct; +import prerna.sablecc2.om.PixelDataType; +import prerna.sablecc2.om.ReactorKeysEnum; +import prerna.sablecc2.om.nounmeta.NounMetadata; +import prerna.util.Utility; + +/** + * List all versions of a specific object in a versioned storage engine. + * + *

Pixel 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()); + + if (!storage.isVersioningEnabled()) { + throw new IllegalArgumentException("Storage engine does not have versioning enabled"); + } + + 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"); + } +} From 1cd59c471e84b629deb831d32b5d852e33e48952 Mon Sep 17 00:00:00 2001 From: Matt Freshwaters Date: Thu, 11 Jun 2026 10:11:36 -0600 Subject: [PATCH 5/6] fix: resolved isLatest detection for GCS object versions --- .../engine/impl/storage/GoogleCloudNativeBlobStorageEngine.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/prerna/engine/impl/storage/GoogleCloudNativeBlobStorageEngine.java b/src/prerna/engine/impl/storage/GoogleCloudNativeBlobStorageEngine.java index ddbb8f2f2d0..804c19c18bb 100644 --- a/src/prerna/engine/impl/storage/GoogleCloudNativeBlobStorageEngine.java +++ b/src/prerna/engine/impl/storage/GoogleCloudNativeBlobStorageEngine.java @@ -276,7 +276,7 @@ public List> listVersions(String storagePath) throws Excepti versionInfo.put("lastModified", blob.getUpdateTimeOffsetDateTime() != null ? blob.getUpdateTimeOffsetDateTime().toString() : null); versionInfo.put("size", blob.getSize()); - versionInfo.put("isLatest", !blob.isDirectory()); + versionInfo.put("isLatest", blob.getDeleteTimeOffsetDateTime() == null); versionInfo.put("key", blob.getName()); versions.add(versionInfo); } From 7234b5de370149333bbba9ef75532ee6e1aaed09 Mon Sep 17 00:00:00 2001 From: "Patel, Parth" Date: Thu, 11 Jun 2026 16:28:45 -0400 Subject: [PATCH 6/6] refactor: refactoring entire flow to not use smss --- src/prerna/engine/api/IStorageEngine.java | 31 +------ .../storage/AWSNativeBlobStorageEngine.java | 87 +----------------- .../storage/AbstractRCloneStorageEngine.java | 3 +- .../impl/storage/AbstractStorageEngine.java | 18 ---- .../storage/AzureNativeBlobStorageEngine.java | 3 +- ...DeveloperLocalFileSystemStorageEngine.java | 3 +- .../GoogleCloudNativeBlobStorageEngine.java | 92 +------------------ .../impl/storage/JCIFSStorageEngine.java | 4 +- .../impl/storage/SFTPStorageEngine.java | 3 +- .../storage/ListStorageVersionsReactor.java | 4 - .../reactor/storage/PushToStorageReactor.java | 17 ++-- 11 files changed, 27 insertions(+), 238 deletions(-) diff --git a/src/prerna/engine/api/IStorageEngine.java b/src/prerna/engine/api/IStorageEngine.java index 62952e74d59..e3c39909642 100644 --- a/src/prerna/engine/api/IStorageEngine.java +++ b/src/prerna/engine/api/IStorageEngine.java @@ -94,39 +94,16 @@ default List> listVersions(String storagePath) throws Except /** * 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 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; - - /** - * Check if object versioning is enabled for this storage engine. - * - * @return true if versioning is enabled, false by default - */ - default boolean isVersioningEnabled() { - return false; - } - - /** - * Copy local files to storage and return the version identifier of the uploaded object. - * Only supported by engines with versioning enabled (S3, GCS). - * Default implementation delegates to {@link #copyToStorage} and returns null. - * - * @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 supported - * @throws Exception if the upload fails - */ - default String copyToStorageVersioned(String localFilePath, String storageFolderPath, - Map metadata) throws Exception { - copyToStorage(localFilePath, storageFolderPath, metadata); - return null; - } + String copyToStorage(String localFilePath, String storageFolderPath, Map metadata) throws Exception; /** * diff --git a/src/prerna/engine/impl/storage/AWSNativeBlobStorageEngine.java b/src/prerna/engine/impl/storage/AWSNativeBlobStorageEngine.java index cdbe227e94b..450a1c8b89e 100644 --- a/src/prerna/engine/impl/storage/AWSNativeBlobStorageEngine.java +++ b/src/prerna/engine/impl/storage/AWSNativeBlobStorageEngine.java @@ -415,59 +415,7 @@ public void syncStorageToLocal(String storagePath, String localPath) throws Exce } @Override - public void copyToStorage(String localFilePath, String storageFolderPath, Map metadata) - throws Exception { - List paths = parseLocalPaths(localFilePath); - List uploadedFiles = new ArrayList<>(); - List failedFiles = new ArrayList<>(); - boolean found = false; - - for (Path path : paths) { - if (!Files.exists(path)) { - classLogger.error("File not found: {}", path); - failedFiles.add(path.toString()); - continue; - } - - deleteEmptyDirectories(path); - - if (Files.isDirectory(path)) { - try (Stream stream = Files.walk(path)) { - stream.filter(Files::isRegularFile).forEach(file -> { - try { - uploadedFiles.add(uploadFile(path, file, storageFolderPath, metadata)); - } catch (Exception e) { - failedFiles.add(file.toString()); - classLogger.error("Failed to upload file: {}", file, e); - rollbackUploads(this.client, failedFiles); - } - }); - found = true; - } - } else { - try { - uploadedFiles.add(uploadFile(path.getParent(), path, storageFolderPath, metadata)); - found = true; - } catch (Exception e) { - failedFiles.add(path.toString()); - classLogger.error("Failed to upload file: {}", path, e); - rollbackUploads(this.client, failedFiles); - } - } - } - // Delete empty blobs from S3 - deleteEmptyBlobsFromS3(storageFolderPath); - if (uploadedFiles.isEmpty()) { - classLogger.info("No files were uploaded."); - } else { - classLogger.info("Successfully uploaded files: {}", uploadedFiles); - } - classLogger.info(found ? "Copy completed successfully for: {}" : "No files found to copy for: {}", - storageFolderPath); - } - - @Override - public String copyToStorageVersioned(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<>(); @@ -488,7 +436,7 @@ public String copyToStorageVersioned(String localFilePath, String storageFolderP try (Stream stream = Files.walk(path)) { stream.filter(Files::isRegularFile).forEach(file -> { try { - String versionId = uploadFileVersioned(path, file, storageFolderPath, metadata); + String versionId = uploadFile(path, file, storageFolderPath, metadata); uploadedFiles.add(file.toString()); if (versionId != null) { lastVersionId.set(versionId); @@ -503,7 +451,7 @@ public String copyToStorageVersioned(String localFilePath, String storageFolderP } } else { try { - String versionId = uploadFileVersioned(path.getParent(), path, storageFolderPath, metadata); + String versionId = uploadFile(path.getParent(), path, storageFolderPath, metadata); uploadedFiles.add(path.toString()); if (versionId != null) { lastVersionId.set(versionId); @@ -827,7 +775,7 @@ private String uploadFileToS3(String storagePath, Path filePath, Path basePath, PutObjectResponse response = this.client.putObject(putRequest, filePath); classLogger.info("Uploaded/Updated file: {}", fileKey); - if (this.versioningEnabled && response.versionId() != null) { + if (response.versionId() != null) { classLogger.info("Version ID for {}: {}", fileKey, response.versionId()); } }, "Uploading file to S3: " + fileKey); @@ -999,39 +947,12 @@ private String uploadFile(Path rootPath, Path file, String storageFolderPath, Ma .collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().toString()))); } - retryOperation(() -> { - this.client.putObject(putBuilder.build(), file); - classLogger.info("Uploaded file to S3: {}", fileKey); - }, "Uploading to S3: " + fileKey); - - return fileKey; - } - - private String uploadFileVersioned(Path rootPath, Path file, String storageFolderPath, - Map metadata) throws IOException { - String normalizedPath = Utility.normalizePath(storageFolderPath).trim(); - if (normalizedPath.startsWith("/")) { - normalizedPath = normalizedPath.substring(1); - } - - String relativePath = Utility.normalizePath(rootPath.relativize(file).toString()).trim(); - String fileKey = normalizedPath.isEmpty() ? relativePath - : (normalizedPath.endsWith("/") ? normalizedPath + relativePath : normalizedPath + "/" + relativePath); - - PutObjectRequest.Builder putBuilder = PutObjectRequest.builder().bucket(this.bucket).key(fileKey); - - if (metadata != null && !metadata.isEmpty()) { - putBuilder.metadata(metadata.entrySet().stream() - .collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().toString()))); - } - AtomicReference versionIdRef = new AtomicReference<>(null); retryOperation(() -> { PutObjectResponse response = this.client.putObject(putBuilder.build(), file); classLogger.info("Uploaded file to S3: {}", fileKey); if (response.versionId() != null) { versionIdRef.set(response.versionId()); - classLogger.info("Version ID for {}: {}", fileKey, response.versionId()); } }, "Uploading to S3: " + fileKey); 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/AbstractStorageEngine.java b/src/prerna/engine/impl/storage/AbstractStorageEngine.java index 9fb04d04c4c..979c77eadc7 100644 --- a/src/prerna/engine/impl/storage/AbstractStorageEngine.java +++ b/src/prerna/engine/impl/storage/AbstractStorageEngine.java @@ -47,10 +47,6 @@ public abstract class AbstractStorageEngine extends AbstractEngine implements IS protected static final Gson GSON = new GsonBuilder().disableHtmlEscaping() .setObjectToNumberStrategy(ToNumberPolicy.LONG_OR_DOUBLE).create(); - public static final String VERSIONING_ENABLED_KEY = "VERSIONING_ENABLED"; - - protected boolean versioningEnabled = false; - /** * Init the general storage values * @@ -60,20 +56,6 @@ public abstract class AbstractStorageEngine extends AbstractEngine implements IS @Override public void open(Properties smssProp) throws Exception { super.open(smssProp); - String versioningStr = smssProp.getProperty(VERSIONING_ENABLED_KEY); - if (versioningStr != null && !versioningStr.isEmpty()) { - this.versioningEnabled = Boolean.parseBoolean(versioningStr); - } - } - - /** - * Check if versioning is enabled for this storage engine. - * - * @return true if versioning is enabled - */ - @Override - public boolean isVersioningEnabled() { - return this.versioningEnabled; } /** 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 metadata) - throws Exception { - List paths = parseLocalPaths(localFilePath); - List uploadedFiles = new ArrayList<>(); - List failedFiles = new ArrayList<>(); - boolean found = false; - for (Path filePath : paths) { - if (!Files.exists(filePath)) { - classLogger.error("File not found: {}", filePath); - failedFiles.add(filePath.toString()); - continue; - } - - // Delete empty directories before upload - deleteEmptyDirectories(filePath); - - if (Files.isDirectory(filePath)) { - try (Stream stream = Files.walk(filePath)) { - stream.filter(Files::isRegularFile).forEach(file -> { - try { - uploadedFiles.add(uploadFileToGCS(filePath, file, storageFolderPath, metadata)); - } catch (Exception e) { - failedFiles.add(file.toString()); - classLogger.error("Failed to upload file: {}", file, e); - rollbackUploads(storage, failedFiles); - } - }); - found = true; - } - } else { - try { - uploadedFiles.add(uploadFileToGCS(filePath.getParent(), filePath, storageFolderPath, metadata)); - found = true; - } catch (Exception e) { - failedFiles.add(filePath.toString()); - classLogger.error("Failed to upload file: {}", filePath, e); - rollbackUploads(storage, failedFiles); - } - } - } - - // Delete empty blobs from GCS - deleteEmptyBlobs(storageFolderPath); - - if (uploadedFiles.isEmpty()) { - classLogger.info("No files were uploaded."); - } else { - classLogger.info("Successfully uploaded files: {}", uploadedFiles); - } - classLogger.info(found ? "Copy completed successfully for: {}" : "No files found to copy for: {}", - storageFolderPath); - } - - @Override - public String copyToStorageVersioned(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<>(); @@ -486,7 +432,7 @@ public String copyToStorageVersioned(String localFilePath, String storageFolderP try (Stream stream = Files.walk(filePath)) { stream.filter(Files::isRegularFile).forEach(file -> { try { - String generation = uploadFileToGCSVersioned(filePath, file, storageFolderPath, metadata); + String generation = uploadFileToGCS(filePath, file, storageFolderPath, metadata); uploadedFiles.add(file.toString()); if (generation != null) { lastVersionId.set(generation); @@ -501,8 +447,7 @@ public String copyToStorageVersioned(String localFilePath, String storageFolderP } } else { try { - String generation = uploadFileToGCSVersioned(filePath.getParent(), filePath, storageFolderPath, - metadata); + String generation = uploadFileToGCS(filePath.getParent(), filePath, storageFolderPath, metadata); uploadedFiles.add(filePath.toString()); if (generation != null) { lastVersionId.set(generation); @@ -996,36 +941,6 @@ private String uploadFileToGCS(Path rootPath, Path file, String storageFolderPat .collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().toString()))); } - retryOperation(() -> { - try { - storage.create(blobInfoBuilder.build(), Files.readAllBytes(file)); - classLogger.info("Uploaded file to GCS: {}", blobName); - } catch (IOException e) { - classLogger.error("Failed to upload file to GCS: {}", blobName, e); - } - }, "Uploading: " + blobName); - - return blobName; - } - - private String uploadFileToGCSVersioned(Path rootPath, Path file, String storageFolderPath, - Map metadata) throws IOException { - String normalizedPath = Utility.normalizePath(storageFolderPath).trim(); - if (normalizedPath.startsWith("/")) { - normalizedPath = normalizedPath.substring(1); - } - String relativePath = Utility.normalizePath(rootPath.relativize(file).toString()).trim(); - String blobName = normalizedPath.isEmpty() ? relativePath - : (normalizedPath.endsWith("/") ? normalizedPath + relativePath : normalizedPath + "/" + relativePath); - - BlobId blobId = BlobId.of(this.BUCKET, blobName); - BlobInfo.Builder blobInfoBuilder = BlobInfo.newBuilder(blobId); - - if (metadata != null && !metadata.isEmpty()) { - blobInfoBuilder.setMetadata(metadata.entrySet().stream() - .collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().toString()))); - } - AtomicReference generationRef = new AtomicReference<>(null); retryOperation(() -> { try { @@ -1033,7 +948,6 @@ private String uploadFileToGCSVersioned(Path rootPath, Path file, String storage classLogger.info("Uploaded file to GCS: {}", blobName); if (blob.getGeneration() != null) { generationRef.set(String.valueOf(blob.getGeneration())); - classLogger.info("Generation for {}: {}", blobName, blob.getGeneration()); } } catch (IOException e) { classLogger.error("Failed to upload file to GCS: {}", blobName, e); 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, Map> versions = storage.listVersions(storagePath); return new NounMetadata(versions, PixelDataType.VECTOR); diff --git a/src/prerna/reactor/storage/PushToStorageReactor.java b/src/prerna/reactor/storage/PushToStorageReactor.java index 05aeb0d65ba..cda2dd30fbc 100644 --- a/src/prerna/reactor/storage/PushToStorageReactor.java +++ b/src/prerna/reactor/storage/PushToStorageReactor.java @@ -86,17 +86,12 @@ public NounMetadata execute() { Map metadata = getMetadata(); try { - // If this engine supports versioning, use the versioned upload - if (storage.isVersioningEnabled()) { - String versionId = storage.copyToStorageVersioned(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); - } - } else { - 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) {