Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,7 @@
<version.slf4j>2.0.17</version.slf4j>
<version.testcontainers>1.21.3</version.testcontainers>
<version.threeten>1.8.0</version.threeten>
<version.azureBlob>12.25.0</version.azureBlob>
</properties>

<dependencyManagement>
Expand Down Expand Up @@ -309,6 +310,12 @@
<version>${version.logback}</version>
<scope>test</scope>
</dependency>

<dependency>
<groupId>com.azure</groupId>
<artifactId>azure-storage-blob</artifactId>
<version>${version.azureBlob}</version>
</dependency>
</dependencies>

<build>
Expand Down
Original file line number Diff line number Diff line change
@@ -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<String, Boolean> fileExistsCache = new ConcurrentHashMap<>();
private final Map<String, List<String>> directoryListCache = new ConcurrentHashMap<>();
private final Map<String, byte[]> smallFileCache = new ConcurrentHashMap<>();
private final Map<String, Long> fileSizeCache = new ConcurrentHashMap<>();
private volatile List<String> 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<BlobCorsRule> 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<String, Boolean> fileExistsBatch(List<String> paths) {
Map<String, Boolean> results = new ConcurrentHashMap<>();

List<CompletableFuture<Void>> 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<String> 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<String> listFilesInternal(String directoryPath) {
return directoryListCache.computeIfAbsent(directoryPath, path -> {
String prefix = fullPath(path);
PagedIterable<BlobItem> 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<String> getAllFiles() throws IOException {
if (allFilesCache != null) return allFilesCache;

synchronized (this) {
if (allFilesCache != null) return allFilesCache;

PagedIterable<BlobItem> items = containerClient.listBlobsByHierarchy(rootPath);
List<String> 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<String> allFiles = getAllFiles();
if (allFiles.isEmpty()) return 0;

long total = 0;
List<CompletableFuture<Long>> 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<Long> f : futures) {
try {
total += f.join();
} catch (Exception e) {

}
}

return total;
}

public Map<String, Integer> 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;
}
}