feat: initial commit
This commit is contained in:
@@ -0,0 +1,336 @@
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.IO;
|
||||
using System.Linq;
|
||||
using System.Text.RegularExpressions;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using Microsoft.Extensions.Options;
|
||||
using Moq;
|
||||
using NUnit.Framework;
|
||||
using Snowflake.Data.Client;
|
||||
using Strata.Analytics.Biz.DataQuery;
|
||||
using Strata.Analytics.Biz.Extracts;
|
||||
using Strata.Analytics.Biz.Extracts.Export;
|
||||
using Strata.Analytics.Biz.Extracts.Jobs;
|
||||
using Strata.Analytics.Biz.Extracts.Models;
|
||||
using Strata.Analytics.Biz.Extracts.Services;
|
||||
using Strata.Analytics.Biz.Hangfire;
|
||||
using Strata.Analytics.Test.Unit.Biz.TestData;
|
||||
using Strata.CoreLib.Claims;
|
||||
using Strata.DataSchema.Models.Query;
|
||||
|
||||
namespace Strata.Analytics.Test.Unit.Biz.Extracts
|
||||
{
|
||||
[TestFixture]
|
||||
internal class ExtractExportJobTest
|
||||
{
|
||||
private Mock<IS3BucketService> _s3BucketService;
|
||||
private Mock<IExtractService> _extractService;
|
||||
private Mock<IExtractExportHistoryService> _historyService;
|
||||
private Mock<IRecurringJobService> _recurringJobService;
|
||||
private Mock<IExtractScheduleService> _scheduleService;
|
||||
private Mock<ISftpService> _sftpService;
|
||||
private Mock<IDataQueryService> _dataQueryService;
|
||||
private Mock<ILogger<ExtractExportJob>> _logger;
|
||||
private IClaimsPrincipalAccessor _claimsPrincipalAccessor;
|
||||
private Mock<IOptions<ExtractConfiguration>> _config;
|
||||
|
||||
|
||||
[Test]
|
||||
public void ExecuteExport_3Files_NoOutbound()
|
||||
{
|
||||
const string SFTP_PATH = "SFTP://directory/folder/subpath";
|
||||
var exportJob = CreateExtractExportJob();
|
||||
var extract = GetRawSqlExtract();
|
||||
var exportRequest = GetExtractExportRequest(extract);
|
||||
var externalStageTokens = new ExternalStageUrlTokens
|
||||
{ BucketName = "Big Bucket", FolderPath = "s3://somewhere/out/there/" };
|
||||
var files = GetFiles(exportRequest.OutputFileName, externalStageTokens, 3);
|
||||
var outboundFiles = files.Select(x => new OutboundFileInfo() { Destination = SFTP_PATH, FileName = x.FilePath, FileSize = x.FileSize, RowCount = x.RowCount }).ToList();
|
||||
|
||||
_extractService.Setup(x => x.GetByIdForExecutionAsync(It.IsAny<int>(), It.IsAny<CancellationToken>())).Returns(Task.FromResult(extract));
|
||||
_dataQueryService.Setup(x => x.BuildSqlQuery(It.IsAny<QueryConfig>(), It.IsAny<CancellationToken>())).Returns(Task.FromResult(new SqlResponse { SqlQuery = extract.RawSql }));
|
||||
_dataQueryService.Setup(x => x.GetExternalStageUrlTokensAsync(It.IsAny<CancellationToken>())).Returns(Task.FromResult(externalStageTokens));
|
||||
_dataQueryService.Setup(x => x.ExecuteSqlQueryAsync(It.IsAny<string>(), It.IsAny<IEnumerable<KeyValuePair<string, object>>>())).Returns(Task.FromResult(new List<dynamic>().AsEnumerable()));
|
||||
_s3BucketService
|
||||
.Setup(x => x.GetFiles(It.IsAny<string>(), It.IsAny<string>(), It.IsAny<bool>(),
|
||||
It.IsAny<string>(), It.IsAny<int>(), It.IsAny<CancellationToken>()))
|
||||
.Returns(Task.FromResult(files));
|
||||
_sftpService.Setup(x => x.UploadFilesFromS3(It.IsAny<string>(), files, It.IsAny<string>(), It.IsAny<string>(), It.IsAny<CancellationToken>())).Returns(Task.FromResult(outboundFiles));
|
||||
_historyService
|
||||
.Setup(x => x.CreateAsync(It.IsAny<ExtractExportHistorySaveInfo>(), It.IsAny<CancellationToken>()))
|
||||
.Returns(Task.FromResult(1));
|
||||
_historyService
|
||||
.Setup(x => x.UpdateAsync(It.IsAny<int>(), It.IsAny<ExtractExportHistorySaveInfo>(), It.IsAny<CancellationToken>()))
|
||||
.Returns(Task.FromResult(new ExtractExportHistory
|
||||
{ ExtractId = extract.ExtractId, ExtractHistoryId = 1, ExportStartDate = DateTimeOffset.UtcNow, ExportEndDate = DateTimeOffset.UtcNow }));
|
||||
_historyService
|
||||
.Setup(x => x.HasRunningExtractByExtractId(It.IsAny<int>(), It.IsAny<CancellationToken>()))
|
||||
.Returns(Task.FromResult(false));
|
||||
|
||||
//assert
|
||||
Assert.DoesNotThrowAsync(() => exportJob.ExecuteExport(It.IsAny<string>(),It.IsAny<Guid>(),It.IsAny<int>(), It.IsAny<Guid>(), exportRequest, new PerformContextMock().Object, It.IsAny<CancellationToken>()));
|
||||
_extractService.Verify(s => s.GetByIdForExecutionAsync(extract.ExtractId, It.IsAny<CancellationToken>()), Times.Once);
|
||||
_dataQueryService.Verify(s => s.GetExternalStageUrlTokensAsync(It.IsAny<CancellationToken>()), Times.Once);
|
||||
_dataQueryService.Verify(s => s.ExecuteSqlQueryAsync(It.Is<string>(query => query.Contains(extract.RawSql)), It.IsAny<IEnumerable<KeyValuePair<string, object>>>()), Times.Once);
|
||||
_s3BucketService.Verify(s => s.GetFiles(It.IsAny<string>(), It.IsAny<string>(), It.IsAny<bool>(),
|
||||
It.IsAny<string>(), It.IsAny<int>(), It.IsAny<CancellationToken>()), Times.Once);
|
||||
_sftpService.Verify(
|
||||
s => s.UploadFilesFromS3(externalStageTokens.BucketName,
|
||||
files, It.IsAny<string>(), It.IsAny<string>(), It.IsAny<CancellationToken>()), Times.Once);
|
||||
_historyService.Verify(
|
||||
s => s.HasRunningExtractByExtractId(It.IsAny<int>(), It.IsAny<CancellationToken>()),
|
||||
Times.Once);
|
||||
_historyService.Verify(
|
||||
s => s.CreateAsync(
|
||||
It.Is<ExtractExportHistorySaveInfo>(h => HistorySaveInfosAreEquivalent(h,
|
||||
new ExtractExportHistorySaveInfo
|
||||
{
|
||||
ExtractId = extract.ExtractId,
|
||||
OutputFilename = exportRequest.OutputFileName,
|
||||
SqlExecuted = extract.RawSql,
|
||||
TotalFileSize = files.Sum(f => f.FileSize),
|
||||
TotalRowCount = files.Sum(f => f.RowCount),
|
||||
Destination = SFTP_PATH
|
||||
})),
|
||||
It.IsAny<CancellationToken>()),
|
||||
Times.Once);
|
||||
}
|
||||
|
||||
private static bool HistorySaveInfosAreEquivalent(ExtractExportHistorySaveInfo actual, ExtractExportHistorySaveInfo expected)
|
||||
{
|
||||
Assert.That(actual.ExtractId, Is.EqualTo(expected.ExtractId));
|
||||
Assert.That(actual.OutputFilename, Is.EqualTo(expected.OutputFilename));
|
||||
Assert.That(actual.SqlExecuted.Contains(expected.SqlExecuted));
|
||||
Assert.That(actual.TotalFileSize, Is.EqualTo(expected.TotalFileSize));
|
||||
Assert.That(actual.TotalRowCount, Is.EqualTo(expected.TotalRowCount));
|
||||
Assert.That(actual.Destination, Is.EqualTo(expected.Destination));
|
||||
Assert.That(actual.ExportStartDate, Is.LessThan(actual.ExportEndDate));
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
[Test]
|
||||
public void ExecuteExport_FileSizeException()
|
||||
{
|
||||
var exportJob = CreateExtractExportJob();
|
||||
var extract = GetRawSqlExtract();
|
||||
var exportRequest = GetExtractExportRequest(extract);
|
||||
var externalStageTokens = new ExternalStageUrlTokens
|
||||
{ BucketName = "Big Bucket", FolderPath = "s3://somewhere/out/there/" };
|
||||
var files = GetFiles(exportRequest.OutputFileName, externalStageTokens, 3);
|
||||
|
||||
_extractService.Setup(x => x.GetByIdForExecutionAsync(It.IsAny<int>(), It.IsAny<CancellationToken>())).Returns(Task.FromResult(extract));
|
||||
_dataQueryService.Setup(x => x.BuildSqlQuery(It.IsAny<QueryConfig>(), It.IsAny<CancellationToken>())).Returns(Task.FromResult(new SqlResponse { SqlQuery = extract.RawSql }));
|
||||
_dataQueryService.Setup(x => x.GetExternalStageUrlTokensAsync(It.IsAny<CancellationToken>())).Returns(Task.FromResult(externalStageTokens));
|
||||
_dataQueryService.Setup(x => x.ExecuteSqlQueryAsync(It.IsAny<string>(), It.IsAny<IEnumerable<KeyValuePair<string, object>>>())).Throws(new SnowflakeDbException(string.Empty, 0, "Max file size exceeded for unload"));
|
||||
_s3BucketService
|
||||
.Setup(x => x.GetFiles(It.IsAny<string>(), It.IsAny<string>(), It.IsAny<bool>(),
|
||||
It.IsAny<string>(), It.IsAny<int>(), It.IsAny<CancellationToken>()))
|
||||
.Returns(Task.FromResult(files));
|
||||
_sftpService.Setup(x => x.UploadFilesFromS3(It.IsAny<string>(), files, "test", It.IsAny<string>(), It.IsAny<CancellationToken>()));
|
||||
_historyService
|
||||
.Setup(x => x.CreateAsync(It.IsAny<ExtractExportHistorySaveInfo>(), It.IsAny<CancellationToken>()))
|
||||
.Returns(Task.FromResult(1));
|
||||
_historyService
|
||||
.Setup(x => x.HasRunningExtractByExtractId(It.IsAny<int>(), It.IsAny<CancellationToken>()))
|
||||
.Returns(Task.FromResult(false));
|
||||
|
||||
Assert.ThrowsAsync<SnowflakeDbException>(() => exportJob.ExecuteExport(It.IsAny<string>(),It.IsAny<Guid>(), It.IsAny<int>(), It.IsAny<Guid>(), exportRequest, new PerformContextMock().Object, It.IsAny<CancellationToken>()));
|
||||
|
||||
_extractService.Verify(s => s.GetByIdForExecutionAsync(extract.ExtractId, It.IsAny<CancellationToken>()), Times.Once);
|
||||
_dataQueryService.Verify(s => s.GetExternalStageUrlTokensAsync(It.IsAny<CancellationToken>()), Times.Once);
|
||||
_dataQueryService.Verify(s => s.ExecuteSqlQueryAsync(It.Is<string>(query => query.Contains(extract.RawSql)), It.IsAny<IEnumerable<KeyValuePair<string, object>>>()), Times.Once);
|
||||
_s3BucketService.Verify(s => s.GetFiles(It.IsAny<string>(), It.IsAny<string>(), It.IsAny<bool>(),
|
||||
It.IsAny<string>(), It.IsAny<int>(), It.IsAny<CancellationToken>()), Times.Never);
|
||||
_sftpService.Verify(
|
||||
s => s.UploadFilesFromS3(externalStageTokens.BucketName,
|
||||
files, "test", It.IsAny<string>(), It.IsAny<CancellationToken>()), Times.Never);
|
||||
|
||||
_historyService.Verify(
|
||||
s => s.HasRunningExtractByExtractId(It.IsAny<int>(), It.IsAny<CancellationToken>()),
|
||||
Times.Once);
|
||||
// should be created
|
||||
_historyService.Verify(
|
||||
s => s.CreateAsync(It.IsAny<ExtractExportHistorySaveInfo>(), It.IsAny<CancellationToken>()),
|
||||
Times.Once);
|
||||
// should be updated with an error status
|
||||
_historyService.Verify(
|
||||
s => s.UpdateAsync(1, It.Is<ExtractExportHistorySaveInfo>(h => h.Status == ExtractExportStatus.Error), It.IsAny<CancellationToken>()),
|
||||
Times.Once);
|
||||
}
|
||||
[Test]
|
||||
public void ExecuteExport_3Files_WithOutbound()
|
||||
{
|
||||
const string SFTP_PATH = "SFTP://directory/folder/subpath";
|
||||
|
||||
var exportJob = CreateExtractExportJob();
|
||||
var extract = GetRawSqlExtract("SELECT * FROM dbo.ALLDATAEVER");
|
||||
var exportRequest = GetExtractExportRequest(extract);
|
||||
var externalStageTokens = new ExternalStageUrlTokens
|
||||
{ BucketName = "Big Bucket", FolderPath = "s3://somewhere/out/there/" };
|
||||
var files = GetFiles(exportRequest.OutputFileName, externalStageTokens, 3);
|
||||
var outboundFiles = files.Select(x => new OutboundFileInfo() { Destination = SFTP_PATH, FileName = x.FilePath, FileSize = x.FileSize, RowCount = x.RowCount }).ToList();
|
||||
|
||||
_extractService.Setup(x => x.GetByIdForExecutionAsync(It.IsAny<int>(), It.IsAny<CancellationToken>())).Returns(Task.FromResult(extract));
|
||||
_dataQueryService.Setup(x => x.GetExternalStageUrlTokensAsync(It.IsAny<CancellationToken>())).Returns(Task.FromResult(externalStageTokens));
|
||||
_dataQueryService.Setup(x => x.ExecuteSqlQueryAsync(It.IsAny<string>(), It.IsAny<IEnumerable<KeyValuePair<string, object>>>())).Returns(Task.FromResult(new List<dynamic>().AsEnumerable()));
|
||||
_s3BucketService
|
||||
.Setup(x => x.GetFiles(It.IsAny<string>(), It.IsAny<string>(), It.IsAny<bool>(),
|
||||
It.IsAny<string>(), It.IsAny<int>(), It.IsAny<CancellationToken>()))
|
||||
.Returns(Task.FromResult(files));
|
||||
_sftpService.Setup(x => x.UploadFilesFromS3(It.IsAny<string>(), files, It.IsAny<string>(), It.IsAny<string>(), It.IsAny<CancellationToken>())).Returns(Task.FromResult(outboundFiles));
|
||||
_historyService
|
||||
.Setup(x => x.CreateAsync(It.IsAny<ExtractExportHistorySaveInfo>(), It.IsAny<CancellationToken>()))
|
||||
.Returns(Task.FromResult(1));
|
||||
_historyService
|
||||
.Setup(x => x.UpdateAsync(It.IsAny<int>(), It.IsAny<ExtractExportHistorySaveInfo>(), It.IsAny<CancellationToken>()))
|
||||
.Returns(Task.FromResult(new ExtractExportHistory
|
||||
{ ExtractId = extract.ExtractId, ExtractHistoryId = 1, ExportStartDate = DateTimeOffset.UtcNow, ExportEndDate = DateTimeOffset.UtcNow }));
|
||||
|
||||
//assert
|
||||
Assert.DoesNotThrowAsync(() => exportJob.ExecuteExport(It.IsAny<string>(),It.IsAny<Guid>(), It.IsAny<int>(), It.IsAny<Guid>(), exportRequest, new PerformContextMock().Object, It.IsAny<CancellationToken>()));
|
||||
_extractService.Verify(s => s.GetByIdForExecutionAsync(extract.ExtractId, It.IsAny<CancellationToken>()), Times.Once);
|
||||
_dataQueryService.Verify(s => s.GetExternalStageUrlTokensAsync(It.IsAny<CancellationToken>()), Times.Once);
|
||||
_dataQueryService.Verify(s => s.ExecuteSqlQueryAsync(It.Is<string>(query => query.Contains(extract.RawSql)), It.IsAny<IEnumerable<KeyValuePair<string, object>>>()), Times.Exactly(2));
|
||||
_s3BucketService.Verify(s => s.GetFiles(It.IsAny<string>(), It.IsAny<string>(), It.IsAny<bool>(),
|
||||
It.IsAny<string>(), It.IsAny<int>(), It.IsAny<CancellationToken>()), Times.Once);
|
||||
_sftpService.Verify(
|
||||
s => s.UploadFilesFromS3(externalStageTokens.BucketName,
|
||||
files, It.IsAny<string>(), It.IsAny<string>(), It.IsAny<CancellationToken>()), Times.Once);
|
||||
_sftpService.Verify(s => s.UploadFilesFromMemory(It.IsAny<MemoryStream>(), It.IsAny<string>(), It.IsAny<CancellationToken>()), Times.Once);
|
||||
_historyService.Verify(
|
||||
s => s.CreateAsync(
|
||||
It.Is<ExtractExportHistorySaveInfo>(h => HistorySaveInfosAreEquivalent(h,
|
||||
new ExtractExportHistorySaveInfo
|
||||
{
|
||||
ExtractId = extract.ExtractId,
|
||||
OutputFilename = exportRequest.OutputFileName,
|
||||
SqlExecuted = extract.RawSql,
|
||||
TotalFileSize = files.Sum(f => f.FileSize),
|
||||
TotalRowCount = files.Sum(f => f.RowCount),
|
||||
Destination = SFTP_PATH
|
||||
})),
|
||||
It.IsAny<CancellationToken>()),
|
||||
Times.Once);
|
||||
}
|
||||
[Test]
|
||||
public async Task TestDataToString()
|
||||
{
|
||||
const string SFTP_PATH = "SFTP://directory/folder/subpath";
|
||||
|
||||
var exportJob = CreateExtractExportJob();
|
||||
var files = new List<OutboundFileInfo>
|
||||
{
|
||||
new()
|
||||
{
|
||||
FileName = "/123123131231/121212412/query_2023-04-07-18-56-36.csv_0_0_0.csv.gz",
|
||||
Destination = SFTP_PATH,
|
||||
FileSize = 100,
|
||||
RowCount = 200
|
||||
},
|
||||
new()
|
||||
{
|
||||
FileName = "/123123131231/121212412/query_2024-04-07-18-56-36.csv_0_0_0.csv.gz",
|
||||
Destination = SFTP_PATH,
|
||||
FileSize = 100,
|
||||
RowCount = 200
|
||||
}
|
||||
};
|
||||
var data = new Dictionary<string, List<string>>()
|
||||
{
|
||||
{ "Column 1", new List<string> {"value1","value2","value3"} },
|
||||
{ "Column 2", new List<string> {"value1","value3"} },
|
||||
{ "Column 3", new List<string> {"value1"} },
|
||||
};
|
||||
|
||||
var stringWriter = await exportJob.CollapseDataAndSendFileAsync(files, data, It.IsAny<DateTime>(), It.IsAny<DateTime>(), It.IsAny<CancellationToken>());
|
||||
var stringToTest = stringWriter.ToString();
|
||||
var numberOfColumns = stringToTest.Count(x => x == '|');
|
||||
var numberOfTimesDataSeparated = stringToTest.Count(x => x == ',');
|
||||
var numberOfNewLines = Regex.Matches(stringToTest, "\r\n").Count;
|
||||
Assert.That(numberOfColumns, Is.EqualTo(24));
|
||||
Assert.That(numberOfTimesDataSeparated, Is.EqualTo(6));
|
||||
Assert.IsFalse(stringToTest.Contains("/123123131231/121212412/"));
|
||||
Assert.That(numberOfNewLines, Is.EqualTo(3));
|
||||
}
|
||||
|
||||
[Test]
|
||||
public void ExecuteExport_ConcurrentExtractExportException()
|
||||
{
|
||||
var exportJob = CreateExtractExportJob();
|
||||
var extract = GetRawSqlExtract();
|
||||
var exportRequest = GetExtractExportRequest(extract);
|
||||
_historyService
|
||||
.Setup(x => x.HasRunningExtractByExtractId(It.IsAny<int>(), It.IsAny<CancellationToken>()))
|
||||
.Returns(Task.FromResult(true));
|
||||
|
||||
Assert.ThrowsAsync<ConcurrentExtractExportException>(() => exportJob.ExecuteExport(It.IsAny<string>(),It.IsAny<Guid>(), It.IsAny<int>(), It.IsAny<Guid>(), exportRequest, new PerformContextMock().Object, It.IsAny<CancellationToken>()));
|
||||
|
||||
_historyService.Verify(
|
||||
s => s.HasRunningExtractByExtractId(It.IsAny<int>(), It.IsAny<CancellationToken>()),
|
||||
Times.Once);
|
||||
|
||||
|
||||
}
|
||||
|
||||
private static Extract GetRawSqlExtract(string outboundFileSql = null)
|
||||
{
|
||||
var extract = new Extract()
|
||||
{
|
||||
ExtractId = 1001,
|
||||
Name = "Extraction Point",
|
||||
Description = "insert something here...",
|
||||
ExtractType = ExtractType.CustomSql,
|
||||
StrataId = 1,
|
||||
RawSql = "SELECT * FROM dbo.ALLDATAEVER",
|
||||
OutboundFileSql = outboundFileSql
|
||||
|
||||
};
|
||||
return extract;
|
||||
}
|
||||
|
||||
private static List<S3FileInfo> GetFiles(string fileName, ExternalStageUrlTokens externalStage, int fileNumber)
|
||||
{
|
||||
var files = new List<S3FileInfo>();
|
||||
for (var i = 0; i < fileNumber; i++)
|
||||
{
|
||||
files.Add(new S3FileInfo { FilePath = $"{externalStage.BucketName}/{externalStage.FolderPath}/{fileName}{i}.csv.gz", FileSize = 100, RowCount = 1000 });
|
||||
}
|
||||
return files;
|
||||
}
|
||||
|
||||
private ExtractExportJob CreateExtractExportJob()
|
||||
{
|
||||
_s3BucketService = new Mock<IS3BucketService>();
|
||||
_extractService = new Mock<IExtractService>();
|
||||
_historyService = new Mock<IExtractExportHistoryService>();
|
||||
_recurringJobService = new Mock<IRecurringJobService>();
|
||||
_sftpService = new Mock<ISftpService>();
|
||||
_dataQueryService = new Mock<IDataQueryService>();
|
||||
_claimsPrincipalAccessor = IClaimsPrincipalAccessorMock.GetClaimsPrincipalAccessor();
|
||||
_logger = new Mock<ILogger<ExtractExportJob>>();
|
||||
_scheduleService = new Mock<IExtractScheduleService>();
|
||||
_config = new Mock<IOptions<ExtractConfiguration>>();
|
||||
|
||||
|
||||
var config = new ExtractConfiguration{FileExtension = "txt", MaxFileSize = 1000, SnowflakeMaxFileSize = 1000, OutboundFileExtension = "csv", MaxRowsPerFile = 10};
|
||||
_config.Setup(c => c.Value).Returns(config);
|
||||
|
||||
return new ExtractExportJob(_s3BucketService.Object, _extractService.Object, _historyService.Object, _recurringJobService.Object, _sftpService.Object,
|
||||
_dataQueryService.Object, _claimsPrincipalAccessor, _logger.Object, _scheduleService.Object, _config.Object);
|
||||
}
|
||||
|
||||
private static ExtractExportRequest GetExtractExportRequest(Extract extract)
|
||||
{
|
||||
return new ExtractExportRequest
|
||||
{
|
||||
ExtractId = extract.ExtractId,
|
||||
IncludeOutboundFile = false,
|
||||
OutputFileName = $"{extract.Name}_extended"
|
||||
};
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user