diff --git a/pom.xml b/pom.xml index 78f9ed75..f5462a1c 100644 --- a/pom.xml +++ b/pom.xml @@ -75,6 +75,7 @@ 2.0.17 1.21.3 1.8.0 + 12.25.0 @@ -309,6 +310,12 @@ ${version.logback} test + + + com.azure + azure-storage-blob + ${version.azureBlob} + diff --git a/src/main/java/dev/jcputney/elearning/parser/impl/access/AzureBlobFileAccess.java b/src/main/java/dev/jcputney/elearning/parser/impl/access/AzureBlobFileAccess.java new file mode 100644 index 00000000..20593734 --- /dev/null +++ b/src/main/java/dev/jcputney/elearning/parser/impl/access/AzureBlobFileAccess.java @@ -0,0 +1,245 @@ +package dev.jcputney.elearning.parser.impl.access; + + +import com.azure.core.http.rest.PagedIterable; +import com.azure.storage.blob.*; +import com.azure.storage.blob.models.*; +import com.azure.storage.blob.sas.BlobContainerSasPermission; +import com.azure.storage.blob.sas.BlobServiceSasSignatureValues; + +import dev.jcputney.elearning.parser.api.FileAccess; + +import java.io.ByteArrayInputStream; +import java.io.IOException; +import java.io.InputStream; +import java.time.OffsetDateTime; +import java.util.*; +import java.util.concurrent.*; +import java.util.stream.Collectors; + +import org.apache.commons.lang3.StringUtils; + + +public class AzureBlobFileAccess implements FileAccess { + + private static final long STREAMING_THRESHOLD = 5 * 1024 * 1024; // 5MB + private static final int MAX_CACHE_SIZE = 1000; + + private final BlobContainerClient containerClient; + private final ExecutorService executorService = Executors.newFixedThreadPool(10); + + private final Map fileExistsCache = new ConcurrentHashMap<>(); + private final Map> directoryListCache = new ConcurrentHashMap<>(); + private final Map smallFileCache = new ConcurrentHashMap<>(); + private final Map fileSizeCache = new ConcurrentHashMap<>(); + private volatile List allFilesCache = null; + + private final String containerSas; + + private final String containerURL; + + private final String rootPath; + + public AzureBlobFileAccess(String connectionString, String containerName, String rootPath, boolean eagerCache, String allowedOrigins) { + BlobServiceClient serviceClient = new BlobServiceClientBuilder().connectionString(connectionString) + .buildClient(); + + this.containerClient = serviceClient.getBlobContainerClient(containerName); + BlobServiceSasSignatureValues values = new BlobServiceSasSignatureValues(OffsetDateTime.now().plusDays(10), + new BlobContainerSasPermission().setReadPermission(true)); + + BlobServiceProperties properties = serviceClient.getProperties(); + List corsRules = new ArrayList<>(); + + allowedOrigins = StringUtils.defaultIfEmpty(allowedOrigins, "*"); + + // if your blob is private you may need to apply cors rules to access the scorm course using cors URL + BlobCorsRule corsRule = new BlobCorsRule() + .setAllowedOrigins(allowedOrigins) + .setAllowedMethods("GET") + .setAllowedHeaders("*") + .setMaxAgeInSeconds(3600); + + corsRules.add(corsRule); + properties.setCors(corsRules); + serviceClient.setProperties(properties); + + containerSas = containerClient.generateSas(values); + containerURL = containerClient.getBlobContainerUrl(); + + this.rootPath = rootPath; + + if (eagerCache) { + try { + getAllFiles(); + } catch (IOException e) { + + } + } + } + + @Override + public boolean fileExistsInternal(String path) { + return fileExistsCache.computeIfAbsent(path, p -> { + String blobPath = fullPath(p); + return containerClient.getBlobClient(blobPath).exists(); + }); + } + + @Override + public Map fileExistsBatch(List paths) { + Map results = new ConcurrentHashMap<>(); + + List> futures = paths.stream() + .map(path -> CompletableFuture.runAsync(() -> { + boolean exists = fileExistsInternal(path); + results.put(path, exists); + }, executorService)) + .toList(); + + CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join(); + return results; + } + + @Override + public void prefetchCommonFiles() { + List commonFiles = List.of("imsmanifest.xml", "cmi5.xml", "xAPI.js", "sendStatement.js", + "manifest.xml", "tincan.xml", "MANIFEST.MF"); + + commonFiles.forEach(file -> { + if (!smallFileCache.containsKey(file) && fileExistsInternal(file)) { + try (InputStream inputStream = getFileContentsInternal(file)) { + byte[] data = inputStream.readAllBytes(); + if (data.length <= STREAMING_THRESHOLD) { + smallFileCache.put(file, data); + } + } catch (IOException e) { + + } + } + }); + } + + @Override + public List listFilesInternal(String directoryPath) { + return directoryListCache.computeIfAbsent(directoryPath, path -> { + String prefix = fullPath(path); + PagedIterable items = containerClient.listBlobsByHierarchy(prefix); + return items.stream() + .filter(blob -> !blob.isPrefix()) + .map(BlobItem::getName) + .map(name -> name.substring(rootPath.length())) // strip rootPath + .collect(Collectors.toList()); + }); + } + + @Override + public InputStream getFileContentsInternal(String path) throws IOException { + byte[] cached = smallFileCache.get(path); + if (cached != null) return new ByteArrayInputStream(cached); + + String blobPath = fullPath(path); + BlobClient blobClient = containerClient.getBlobClient(blobPath); + long fileSize = blobClient.getProperties().getBlobSize(); + fileSizeCache.put(path, fileSize); + + if (fileSize <= STREAMING_THRESHOLD) { + byte[] content = blobClient.downloadContent().toBytes(); + if (smallFileCache.size() >= MAX_CACHE_SIZE) { + smallFileCache.keySet().iterator().remove(); + } + smallFileCache.put(path, content); + return new ByteArrayInputStream(content); + } else { + return blobClient.openInputStream(); + } + } + + @Override + public List getAllFiles() throws IOException { + if (allFilesCache != null) return allFilesCache; + + synchronized (this) { + if (allFilesCache != null) return allFilesCache; + + PagedIterable items = containerClient.listBlobsByHierarchy(rootPath); + List allFiles = items.stream() + .filter(blob -> !blob.isPrefix()) + .map(BlobItem::getName) + .map(name -> name.substring(rootPath.length())) + .collect(Collectors.toList()); + + allFiles.forEach(f -> fileExistsCache.put(f, true)); + allFilesCache = allFiles; + return allFiles; + } + } + + @Override + public void clearCaches() { + fileExistsCache.clear(); + directoryListCache.clear(); + smallFileCache.clear(); + fileSizeCache.clear(); + allFilesCache = null; + } + + @Override + public long getTotalSize() throws IOException { + List allFiles = getAllFiles(); + if (allFiles.isEmpty()) return 0; + + long total = 0; + List> futures = new ArrayList<>(); + + for (String file : allFiles) { + Long cachedSize = fileSizeCache.get(file); + if (cachedSize != null) { + total += cachedSize; + } else { + futures.add(CompletableFuture.supplyAsync(() -> { + BlobClient blobClient = containerClient.getBlobClient(fullPath(file)); + long size = blobClient.getProperties().getBlobSize(); + fileSizeCache.put(file, size); + return size; + }, executorService)); + } + } + + for (CompletableFuture f : futures) { + try { + total += f.join(); + } catch (Exception e) { + + } + } + + return total; + } + + public Map getCacheStats() { + return Map.of( + "fileExistsCache", fileExistsCache.size(), + "directoryListCache", directoryListCache.size(), + "smallFileCache", smallFileCache.size(), + "fileSizeCache", fileSizeCache.size() + ); + } + + public void shutdown() { + executorService.shutdown(); + } + + @Override + public String getRootPath() { + return rootPath; + } + + public String getContainerSas() { + return containerSas; + } + + public String getContainerURL() { + return containerURL; + } +} \ No newline at end of file