Files
2026-06-23 11:18:54 -05:00

336 lines
20 KiB
C#

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"
};
}
}
}