import { DefaultAzureCredential } from "@azure/identity"; import _ from "lodash"; import { hashAPIPath } from "."; const { ContainerClient, BlockBlobClient, BlobServiceClient, BlobSASPermissions, ContainerSASPermissions, generateBlobSASQueryParameters, SASProtocol, } = require("@azure/storage-blob"); const { QueueServiceClient, AccountSASResourceTypes, AccountSASServices, AccountSASPermissions, QueueSASPermissions, QueueSASSignatureValues, generateAccountSASQueryParameters, generateQueueSASQueryParameters, StorageSharedKeyCredential, QueueClient, } = require("@azure/storage-queue"); const STORAGE_PATH = process.env.AZURE_PEDW_STORAGE_ENDPOINT; const STORAGE_CONTAINER = process.env.AZURE_PEDW_CONTAINER; const QUEUE_PATH = process.env.AZURE_PEDW_QUEUE_ENDPOINT; const accountName = process.env.AZURE_STORAGE_ACCOUNT_NAME; export const createContainerSas = async (containerName) => { // Get environment variables // Best practice: create time limits const TEN_MINUTES = 10 * 60 * 1000; const NOW = new Date(); // Best practice: set start time a little before current time to // make sure any clock issues are avoided const TEN_MINUTES_BEFORE_NOW = new Date(NOW.valueOf() - TEN_MINUTES); const TEN_MINUTES_AFTER_NOW = new Date(NOW.valueOf() + TEN_MINUTES); // Best practice: use managed identity - DefaultAzureCredential const blobServiceClient = new BlobServiceClient( `${STORAGE_PATH}`, new DefaultAzureCredential() ); // Best practice: delegation key is time-limited // When using a user delegation key, container must already exist const userDelegationKey = await blobServiceClient.getUserDelegationKey( TEN_MINUTES_BEFORE_NOW, TEN_MINUTES_AFTER_NOW ); // Need only list permission to list blobs const containerPermissions = "rcwltdf"; // Best practice: SAS options are time-limited const sasOptions = { containerName, permissions: ContainerSASPermissions.parse(containerPermissions), protocol: SASProtocol.HttpsAndHttp, startsOn: TEN_MINUTES_BEFORE_NOW, expiresOn: TEN_MINUTES_AFTER_NOW, }; //conLogJSON.stringify(sasOptions)); const sasToken = generateBlobSASQueryParameters( sasOptions, userDelegationKey, accountName ).toString(); return sasToken; }; export const createBlobSas = async (containerName, blobName) => { // Get environment variables const accountName = process.env.AZURE_STORAGE_ACCOUNT_NAME; // Best practice: create time limits const TEN_MINUTES = 10 * 60 * 1000; const NOW = new Date(); // Best practice: set start time a little before current time to // make sure any clock issues are avoided const TEN_MINUTES_BEFORE_NOW = new Date(NOW.valueOf() - TEN_MINUTES); const TEN_MINUTES_AFTER_NOW = new Date(NOW.valueOf() + TEN_MINUTES); // Best practice: use managed identity - DefaultAzureCredential const blobServiceClient = new BlobServiceClient( `https://${accountName}.blob.core.windows.net`, new DefaultAzureCredential() ); // Best practice: delegation key is time-limited // When using a user delegation key, container must already exist const userDelegationKey = await blobServiceClient.getUserDelegationKey( TEN_MINUTES_BEFORE_NOW, TEN_MINUTES_AFTER_NOW ); // Need only create/write permission to upload file const blobPermissionsForAnonymousUser = "rcwt"; // Best practice: SAS options are time-limited const sasOptions = { blobName, containerName, permissions: BlobSASPermissions.parse(blobPermissionsForAnonymousUser), protocol: SASProtocol.HttpsAndHttp, startsOn: TEN_MINUTES_BEFORE_NOW, expiresOn: TEN_MINUTES_AFTER_NOW, }; const sasToken = generateBlobSASQueryParameters( sasOptions, userDelegationKey, accountName ).toString(); return sasToken; }; export const createContainer = async (containerName) => { const creds = new DefaultAzureCredential(); //containerName = containerName.toLowerCase(); // console.log( // "container name:", // containerName, // STORAGE_PATH + "/" + containerName // ); const containerClient = new ContainerClient( `${STORAGE_PATH}/${containerName}`, creds ); const blobServiceClient = new BlobServiceClient(`${STORAGE_PATH}`, creds); const createContainerResponse = await containerClient.createIfNotExists(); console.log( "\n//////////////////\n container name :", containerName, "\n//////////////////\n" ); return containerName; }; export const getContainers = async () => { const creds = new DefaultAzureCredential(); const blobServiceClient = new BlobServiceClient(`${STORAGE_PATH}`, creds); console.log("Containers:"); for await (const container of blobServiceClient.listContainers()) { console.log(`- ${container.name}`); } }; export const getBlobs = async (containerName, casefolderID) => { const containerToken = await createContainerSas(containerName); const sasUrl = `${STORAGE_PATH}/${containerName}?${containerToken}`; const containerClient = new ContainerClient(sasUrl); //conLog"getBlobs " + sasUrl); const blobObj = []; for await (const blob of containerClient.listBlobsFlat({ prefix: casefolderID + "/files/", })) { let blobDocumentType = blob.name .split("/")[2] .slice(0, blob.name.split("/")[2].indexOf("_")); blobObj.push({ "name": blob.name.split("/")[2], "path": blob.name, "documentType": blobDocumentType, "versionId": blob.versionId, "caseObj": casefolderID + "/" + casefolderID + "_case.json", "isCurrentVersion": blob.isCurrentVersion, "contentLength": blob.properties.contentLength, "contentType": blob.contentType, "lastModified": blob.properties.lastModified, "filepath": "/api/file/downloadblob?container=" + containerName + "&casefolderID=" + casefolderID + "&blobname=" + blob.name.split("/")[2], "hashedfilepath": hashAPIPath( "/api/file/downloadblob?container=" + containerName + "&casefolderID=" + casefolderID + "&blobname=" + blob.name.split("/")[2] ), "deletepath": "/api/file/deleteblob?container=" + containerName + "&casefolderID=" + casefolderID + "&blobname=" + blob.name.split("/")[2], "hasheddeletepath": hashAPIPath( "/api/file/deleteblob?container=" + containerName + "&casefolderID=" + casefolderID + "&blobname=" + blob.name.split("/")[2] ), "hashgetblobs": hashAPIPath( "/api/file/getbloblist?container=" + containerName + "&casefolderID=" + casefolderID ), }); } //console.log("blobObj:", blobObj); return blobObj; }; export const createBlob = async (formContent, containerName, caseref) => { const containerToken = await createContainerSas(containerName); const sasUrl = `${STORAGE_PATH}/${containerName}?${containerToken}`; const containerClient = new ContainerClient(sasUrl); formContent = JSON.parse(formContent); let caseID = ""; caseID = _.has(formContent, "pinswg_name") ? formContent.pinswg_name : caseref; const content = JSON.stringify(formContent); console.log(content); const blobName = caseID + "/" + caseID + "_appeal.json"; console.log("blobName:", blobName); const blockBlobClient = containerClient.getBlockBlobClient(blobName); const uploadBlobResponse = await blockBlobClient.upload( content, Buffer.byteLength(content) ); const tags = { containerid: containerName, caseID: caseID, }; console.log("the tags:", tags); const withTags = await blockBlobClient.setTags(tags); const withMeta = await blockBlobClient.setMetadata(tags); return formContent.pinswg_name; }; export const createRepBlob = async (formContent, containerName, caseref) => { const containerToken = await createContainerSas(containerName); const sasUrl = `${STORAGE_PATH}/${containerName}?${containerToken}`; const containerClient = new ContainerClient(sasUrl); formContent = JSON.parse(formContent); let caseID = ""; caseID = _.has(formContent, "pinswg_name") ? formContent.pinswg_name : caseref; const content = JSON.stringify(formContent); console.log(content); const blobName = caseref + "/" + caseID + "_rep.json"; console.log("blobName:", blobName); const blockBlobClient = containerClient.getBlockBlobClient(blobName); const uploadBlobResponse = await blockBlobClient.upload( content, Buffer.byteLength(content) ); const tags = { containerid: containerName, caseID: caseID, blobType: "Representation", }; console.log("the tags:", tags); const withTags = await blockBlobClient.setTags(tags); const withMeta = await blockBlobClient.setMetadata(tags); return formContent.pinswg_name; }; export const deleteBlob = async (containerName, blobName) => { const creds = new DefaultAzureCredential(); const options = { deleteSnapshots: "include", // or 'only' }; const containerToken = await createContainerSas(containerName); const sasUrl = `${STORAGE_PATH}/${containerName}?${containerToken}`; const containerClient = new ContainerClient(sasUrl); const blockBlobClient = containerClient.getBlockBlobClient(blobName); await blockBlobClient.delete(options); console.log(`deleted blob ${blobName}`); return { "deleted": blobName }; }; export const deleteBlobCase = async (containerName, blobName) => { const creds = new DefaultAzureCredential(); const options = { deleteSnapshots: "include", // or 'only' }; const containerToken = await createContainerSas(containerName); const sasUrl = `${STORAGE_PATH}/${containerName}?${containerToken}`; const containerClient = new ContainerClient(sasUrl); //conLog"deleteBlobCase "); console.log("blob to delete:", blobName); for await (const blob of containerClient.listBlobsFlat({ prefix: blobName, })) { console.log(" ------ :", blob.name); containerClient.deleteBlob(blob.name); } containerClient.deleteBlob(blobName); console.log(`deleted blob ${blobName}`); for await (const blob of containerClient.listBlobsFlat()) { console.log(" ------ :", blob.name); } return { "deleted": blobName }; }; export const uploadFile = async (formContent, containerName, foldername) => { const containerToken = await createContainerSas(containerName); const sasUrl = `${STORAGE_PATH}/${containerName}?${containerToken}`; const containerClient = new ContainerClient(sasUrl); const files = formContent; //console.log(...formContent); console.log("files...", files, Object.keys(files).length, foldername); for (const prop in files) { console.log(`files[${prop}] = ${files[prop][0].size}`); const blobName = foldername + "/files/" + files[prop][0].fieldName; console.log(blobName); const blockBlobClient = containerClient.getBlockBlobClient(blobName); const uploadBlobResponse = await blockBlobClient.uploadFile( files[prop][0].path, files[prop].size ); const tags = { containerid: containerName, caseID: foldername, documentType: files[prop][0].fieldName.slice( 0, files[prop][0].fieldName.indexOf("_") ), }; const withTags = await blockBlobClient.setTags(tags); const withMeta = await blockBlobClient.setMetadata(tags); console.log( `Uploaded block blob ${files[prop][0].fieldName} successfully`, uploadBlobResponse.requestId ); } }; export const downloadFile = async (containerName, blobName) => { //console.log.apply(containerName, blobName); const containerToken = await createContainerSas(containerName); const sasUrl = `${STORAGE_PATH}/${containerName}?${containerToken}`; const containerClient = new ContainerClient(sasUrl); const blobClient = containerClient.getBlobClient(blobName); const downloadedBlob = await blobClient.download(0); const downloaded = await streamToBuffer(downloadedBlob.readableStreamBody); return downloaded; }; export const downloadProgressFile = async ( containerName, blobName, casefolderID ) => { //console.log.apply(containerName, blobName); const containerToken = await createContainerSas(containerName); const sasUrl = `${STORAGE_PATH}/${containerName}?${containerToken}`; const containerClient = new ContainerClient(sasUrl); const blobClient = containerClient.getBlobClient(blobName); const downloadedBlob = await blobClient.download(); const downloaded = await streamToBuffer(downloadedBlob.readableStreamBody); console.log("Downloaded blob content:", downloaded.toString()); return JSON.parse(downloaded.toString()); }; export const downloadAllProgressFiles = async ( containerName, progressBlobObj ) => { // console.log( // "/////////////////////////\n downloading files: " + // JSON.stringify(progressBlobObj) + // "\n/////////////////////////\n" // ); const containerToken = await createContainerSas(containerName); const sasUrl = `${STORAGE_PATH}/${containerName}?${containerToken}`; const containerClient = new ContainerClient(sasUrl); let blobClient = {}; let downloadedBlob = {}; let downloaded = ""; let caseBlob = {}; let caseDownloaded = ""; let blobCount = 0; for (const prop in progressBlobObj) { //console.log(`Progress - ${prop}: ${progressBlobObj[prop].path}`); blobClient = containerClient.getBlobClient(progressBlobObj[prop].path); downloadedBlob = await blobClient.download(0); downloaded = downloaded + (await streamToBuffer(downloadedBlob.readableStreamBody)) + ","; blobCount++; } downloaded = downloaded.substring(0, downloaded.length - 1); for (const prop in progressBlobObj) { //console.log(`Cases - ${prop}: ${progressBlobObj[prop].caseObj}`); blobClient = containerClient.getBlobClient( progressBlobObj[prop].caseObj ); caseBlob = await blobClient.download(0); caseDownloaded = caseDownloaded + (await streamToBuffer(caseBlob.readableStreamBody)) + ","; } caseDownloaded = caseDownloaded.substring(0, caseDownloaded.length - 1); return JSON.parse( '{ "@odata.count": ' + blobCount + ',"value": [' + downloaded + '],"case":[' + caseDownloaded + "]}" ); }; export const downloadAllRepsFiles = async (containerName, repsBlobObj) => { // console.log( // "/////////////////////////\n downloading files: " + // JSON.stringify(repsBlobObj) + // "\n/////////////////////////\n" // ); const containerToken = await createContainerSas(containerName); const sasUrl = `${STORAGE_PATH}/${containerName}?${containerToken}`; const containerClient = new ContainerClient(sasUrl); let blobClient = {}; let downloadedBlob = {}; let repArr = []; let downloaded = ""; let caseBlob = {}; let caseDownloaded = ""; let blobCount = 0; for (const prop in repsBlobObj) { //console.log(`Progress - ${prop}: ${repsBlobObj[prop].name}`); blobClient = containerClient.getBlobClient(repsBlobObj[prop].name); downloadedBlob = await blobClient.download(0); let repBlobObj = await streamToBuffer( downloadedBlob.readableStreamBody ); repBlobObj = JSON.parse(repBlobObj); repBlobObj["casereference"] = repsBlobObj[prop].name.split("/")[0]; repBlobObj = JSON.stringify(repBlobObj); downloaded = downloaded + repBlobObj + ","; blobCount++; } downloaded = downloaded.substring(0, downloaded.length - 1); return JSON.parse( '{ "@odata.count": ' + blobCount + ',"value": [' + downloaded + "]}" ); }; const streamToBuffer = async (readableStream) => { return new Promise((resolve, reject) => { const chunks = []; readableStream.on("data", (data) => { chunks.push(data instanceof Buffer ? data : Buffer.from(data)); }); readableStream.on("end", () => { resolve(Buffer.concat(chunks)); }); readableStream.on("error", reject); }); }; export const getCaseBlob = async ( containerName, caseReference, formContent ) => { const content = JSON.stringify(formContent); const blobName = caseReference + "/case/" + caseReference + "_case.json"; const containerBlobToken = await createBlobSas(containerName, blobName); const blobSasUrl = `${STORAGE_PATH}/${containerName}/${blobName}?${containerBlobToken}`; const blockBlobClient = new BlockBlobClient(blobSasUrl); const uploadBlobResponse = await blockBlobClient.upload( content, Buffer.byteLength(content) ); return blobName; }; export const getProgressBlobs = async (containerName, caseReference) => { const containerToken = await createContainerSas(containerName); const sasUrl = `${STORAGE_PATH}/${containerName}?${containerToken}`; const containerClient = new ContainerClient(sasUrl); let blobCount = 0; for await (const blob of containerClient.listBlobsFlat({ prefix: caseReference, })) { blobCount++; } console.log( "this is the casefolder and blobcount:", caseReference, blobCount ); let blobObj = []; for await (const blob of containerClient.listBlobsFlat({ prefix: caseReference + "/" + caseReference + "_appeal.json", })) { console.log("getProgressBlobs in here"); blobObj.push({ "name": blob.name.split("/")[1], "path": blob.name, "versionId": blob.versionId, "caseObj": caseReference + "/" + caseReference + "_case.json", "isCurrentVersion": blob.isCurrentVersion, "contentLength": blob.properties.contentLength, "contentType": blob.contentType, "lastModified": blob.properties.lastModified, "hashedfilepath": hashAPIPath( "/api/file/downloadblob?container=" + containerName + "&casefolderID=" + caseReference + "&blobname=" + blob.name.split("/")[1] ), "hasheddeletepath": hashAPIPath( "/api/file/deleteblob?container=" + containerName + "&casefolderID=" + caseReference + "&blobname=" + blob.name.split("/")[1] ), "hashgetblobs": hashAPIPath( "/api/file/getbloblist?container=" + containerName + "&casefolderID=" + caseReference ), }); } blobObj = _.sortBy(blobObj, [ function (o) { return o.lastModified; }, ]).reverse()[0]; return blobObj; }; export const getAllProgressBlobs = async (containerName) => { const containerToken = await createContainerSas(containerName); const sasUrl = `${STORAGE_PATH}/${containerName}?${containerToken}`; const containerClient = new ContainerClient(sasUrl); let blobCount = 0; let blobObj = []; for await (const blob of containerClient.listBlobsFlat()) { blob.name.split("/")[1].indexOf("_appeal.json") > 0 && blob.name.split("/")[1].indexOf("undefined") < 0 && blobObj.push({ "name": blob.name.split("/")[1], "path": blob.name, "versionId": blob.versionId, "caseObj": blob.name.split("/")[0] + "/case/" + blob.name.split("/")[0] + "_case.json", "isCurrentVersion": blob.isCurrentVersion, "contentLength": blob.properties.contentLength, "contentType": blob.contentType, "lastModified": blob.properties.lastModified, }); } //console.log("blobObj:", blobObj); return blobObj; }; export const getRepslobs = async (containerName, caseReference) => { const containerToken = await createContainerSas(containerName); const sasUrl = `${STORAGE_PATH}/${containerName}?${containerToken}`; const containerClient = new ContainerClient(sasUrl); let blobCount = 0; for await (const blob of containerClient.findBlobsByTags( "blobType='Representation'" )) { blobCount++; } console.log( "this is the casefolder and blobcount:", caseReference, blobCount ); let blobObj = []; for await (const blob of containerClient.findBlobsByTags( "blobType='Representation'" )) { // console.log("getRepslobs in here"); // console.log(blob); blobObj.push(blob); } //console.log("blobObjwwwww:", blobObj); return blobObj; }; export const createQueueSas = async (queueName) => { // Get environment variables const account = process.env.AZURE_STORAGE_ACCOUNT_NAME; const accountKey = process.env.AZURE_STORAGE_ACCOUNT_KEY; const sharedKeyCredential = new StorageSharedKeyCredential( account, accountKey ); // Best practice: create time limits const TEN_MINUTES = 10 * 60 * 1000; const NOW = new Date(); // Best practice: set start time a little before current time to // make sure any clock issues are avoided const TEN_MINUTES_AFTER_NOW = new Date(NOW.valueOf() + TEN_MINUTES); var resource_types = new AccountSASResourceTypes(); resource_types.container = true; resource_types.object = true; resource_types.service = true; var permission = new QueueSASPermissions(); permission.read = true; permission.write = true; permission.delete = true; permission.list = true; permission.add = true; permission.update = true; permission.process = true; var sas_signature_val = { expiresOn: TEN_MINUTES_AFTER_NOW, permissions: permission, queueName: queueName, }; const sasToken = generateQueueSASQueryParameters( sas_signature_val, sharedKeyCredential ); //console.log("sass token:", sasToken); return sasToken.toString(); }; export const createCaseCompleteMessage = async ( containerName, caseReference ) => { const whichQueue = "pedw-submitted-applications"; const queueToken = await createQueueSas(whichQueue); const sasUrl = `${QUEUE_PATH}?${queueToken}`; const queueServiceClient = new QueueServiceClient(sasUrl); const message = { containerName: containerName, casepath: containerName + "/" + caseReference + "/" + caseReference + "_appeal.json", casecontentpath: containerName + "/" + caseReference + "/case", filespath: containerName + "/" + caseReference + "/files", // uploadUrl: fileRecord.uploadUrl, // filename: fileRecord.file.name, // fileSize: fileRecord.file.size, }; const sendMessageResponse = await queueServiceClient .getQueueClient(whichQueue) .sendMessage( JSON.stringify( `${JSON.stringify( message )}` ) ); return sendMessageResponse; }; export const createRepCompleteMessage = async ( containerName, caseReference ) => { const whichQueue = "pedw-submitted-representations"; const queueToken = await createQueueSas(whichQueue); const sasUrl = `${QUEUE_PATH}?${queueToken}`; const queueServiceClient = new QueueServiceClient(sasUrl); const message = { containerName: containerName, caseref: caseReference, reppath: containerName + "/" + caseReference + "/", filespath: containerName + "/" + caseReference + "/files", }; const sendMessageResponse = await queueServiceClient .getQueueClient(whichQueue) .sendMessage( JSON.stringify( `${JSON.stringify( message )}` ) ); return sendMessageResponse; };