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 _s3BucketService; private Mock _extractService; private Mock _historyService; private Mock _recurringJobService; private Mock _scheduleService; private Mock _sftpService; private Mock _dataQueryService; private Mock> _logger; private IClaimsPrincipalAccessor _claimsPrincipalAccessor; private Mock> _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(), It.IsAny())).Returns(Task.FromResult(extract)); _dataQueryService.Setup(x => x.BuildSqlQuery(It.IsAny(), It.IsAny())).Returns(Task.FromResult(new SqlResponse { SqlQuery = extract.RawSql })); _dataQueryService.Setup(x => x.GetExternalStageUrlTokensAsync(It.IsAny())).Returns(Task.FromResult(externalStageTokens)); _dataQueryService.Setup(x => x.ExecuteSqlQueryAsync(It.IsAny(), It.IsAny>>())).Returns(Task.FromResult(new List().AsEnumerable())); _s3BucketService .Setup(x => x.GetFiles(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) .Returns(Task.FromResult(files)); _sftpService.Setup(x => x.UploadFilesFromS3(It.IsAny(), files, It.IsAny(), It.IsAny(), It.IsAny())).Returns(Task.FromResult(outboundFiles)); _historyService .Setup(x => x.CreateAsync(It.IsAny(), It.IsAny())) .Returns(Task.FromResult(1)); _historyService .Setup(x => x.UpdateAsync(It.IsAny(), It.IsAny(), It.IsAny())) .Returns(Task.FromResult(new ExtractExportHistory { ExtractId = extract.ExtractId, ExtractHistoryId = 1, ExportStartDate = DateTimeOffset.UtcNow, ExportEndDate = DateTimeOffset.UtcNow })); _historyService .Setup(x => x.HasRunningExtractByExtractId(It.IsAny(), It.IsAny())) .Returns(Task.FromResult(false)); //assert Assert.DoesNotThrowAsync(() => exportJob.ExecuteExport(It.IsAny(),It.IsAny(),It.IsAny(), It.IsAny(), exportRequest, new PerformContextMock().Object, It.IsAny())); _extractService.Verify(s => s.GetByIdForExecutionAsync(extract.ExtractId, It.IsAny()), Times.Once); _dataQueryService.Verify(s => s.GetExternalStageUrlTokensAsync(It.IsAny()), Times.Once); _dataQueryService.Verify(s => s.ExecuteSqlQueryAsync(It.Is(query => query.Contains(extract.RawSql)), It.IsAny>>()), Times.Once); _s3BucketService.Verify(s => s.GetFiles(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny()), Times.Once); _sftpService.Verify( s => s.UploadFilesFromS3(externalStageTokens.BucketName, files, It.IsAny(), It.IsAny(), It.IsAny()), Times.Once); _historyService.Verify( s => s.HasRunningExtractByExtractId(It.IsAny(), It.IsAny()), Times.Once); _historyService.Verify( s => s.CreateAsync( It.Is(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()), 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(), It.IsAny())).Returns(Task.FromResult(extract)); _dataQueryService.Setup(x => x.BuildSqlQuery(It.IsAny(), It.IsAny())).Returns(Task.FromResult(new SqlResponse { SqlQuery = extract.RawSql })); _dataQueryService.Setup(x => x.GetExternalStageUrlTokensAsync(It.IsAny())).Returns(Task.FromResult(externalStageTokens)); _dataQueryService.Setup(x => x.ExecuteSqlQueryAsync(It.IsAny(), It.IsAny>>())).Throws(new SnowflakeDbException(string.Empty, 0, "Max file size exceeded for unload")); _s3BucketService .Setup(x => x.GetFiles(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) .Returns(Task.FromResult(files)); _sftpService.Setup(x => x.UploadFilesFromS3(It.IsAny(), files, "test", It.IsAny(), It.IsAny())); _historyService .Setup(x => x.CreateAsync(It.IsAny(), It.IsAny())) .Returns(Task.FromResult(1)); _historyService .Setup(x => x.HasRunningExtractByExtractId(It.IsAny(), It.IsAny())) .Returns(Task.FromResult(false)); Assert.ThrowsAsync(() => exportJob.ExecuteExport(It.IsAny(),It.IsAny(), It.IsAny(), It.IsAny(), exportRequest, new PerformContextMock().Object, It.IsAny())); _extractService.Verify(s => s.GetByIdForExecutionAsync(extract.ExtractId, It.IsAny()), Times.Once); _dataQueryService.Verify(s => s.GetExternalStageUrlTokensAsync(It.IsAny()), Times.Once); _dataQueryService.Verify(s => s.ExecuteSqlQueryAsync(It.Is(query => query.Contains(extract.RawSql)), It.IsAny>>()), Times.Once); _s3BucketService.Verify(s => s.GetFiles(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny()), Times.Never); _sftpService.Verify( s => s.UploadFilesFromS3(externalStageTokens.BucketName, files, "test", It.IsAny(), It.IsAny()), Times.Never); _historyService.Verify( s => s.HasRunningExtractByExtractId(It.IsAny(), It.IsAny()), Times.Once); // should be created _historyService.Verify( s => s.CreateAsync(It.IsAny(), It.IsAny()), Times.Once); // should be updated with an error status _historyService.Verify( s => s.UpdateAsync(1, It.Is(h => h.Status == ExtractExportStatus.Error), It.IsAny()), 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(), It.IsAny())).Returns(Task.FromResult(extract)); _dataQueryService.Setup(x => x.GetExternalStageUrlTokensAsync(It.IsAny())).Returns(Task.FromResult(externalStageTokens)); _dataQueryService.Setup(x => x.ExecuteSqlQueryAsync(It.IsAny(), It.IsAny>>())).Returns(Task.FromResult(new List().AsEnumerable())); _s3BucketService .Setup(x => x.GetFiles(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) .Returns(Task.FromResult(files)); _sftpService.Setup(x => x.UploadFilesFromS3(It.IsAny(), files, It.IsAny(), It.IsAny(), It.IsAny())).Returns(Task.FromResult(outboundFiles)); _historyService .Setup(x => x.CreateAsync(It.IsAny(), It.IsAny())) .Returns(Task.FromResult(1)); _historyService .Setup(x => x.UpdateAsync(It.IsAny(), It.IsAny(), It.IsAny())) .Returns(Task.FromResult(new ExtractExportHistory { ExtractId = extract.ExtractId, ExtractHistoryId = 1, ExportStartDate = DateTimeOffset.UtcNow, ExportEndDate = DateTimeOffset.UtcNow })); //assert Assert.DoesNotThrowAsync(() => exportJob.ExecuteExport(It.IsAny(),It.IsAny(), It.IsAny(), It.IsAny(), exportRequest, new PerformContextMock().Object, It.IsAny())); _extractService.Verify(s => s.GetByIdForExecutionAsync(extract.ExtractId, It.IsAny()), Times.Once); _dataQueryService.Verify(s => s.GetExternalStageUrlTokensAsync(It.IsAny()), Times.Once); _dataQueryService.Verify(s => s.ExecuteSqlQueryAsync(It.Is(query => query.Contains(extract.RawSql)), It.IsAny>>()), Times.Exactly(2)); _s3BucketService.Verify(s => s.GetFiles(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny()), Times.Once); _sftpService.Verify( s => s.UploadFilesFromS3(externalStageTokens.BucketName, files, It.IsAny(), It.IsAny(), It.IsAny()), Times.Once); _sftpService.Verify(s => s.UploadFilesFromMemory(It.IsAny(), It.IsAny(), It.IsAny()), Times.Once); _historyService.Verify( s => s.CreateAsync( It.Is(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()), Times.Once); } [Test] public async Task TestDataToString() { const string SFTP_PATH = "SFTP://directory/folder/subpath"; var exportJob = CreateExtractExportJob(); var files = new List { 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>() { { "Column 1", new List {"value1","value2","value3"} }, { "Column 2", new List {"value1","value3"} }, { "Column 3", new List {"value1"} }, }; var stringWriter = await exportJob.CollapseDataAndSendFileAsync(files, data, It.IsAny(), It.IsAny(), It.IsAny()); 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(), It.IsAny())) .Returns(Task.FromResult(true)); Assert.ThrowsAsync(() => exportJob.ExecuteExport(It.IsAny(),It.IsAny(), It.IsAny(), It.IsAny(), exportRequest, new PerformContextMock().Object, It.IsAny())); _historyService.Verify( s => s.HasRunningExtractByExtractId(It.IsAny(), It.IsAny()), 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 GetFiles(string fileName, ExternalStageUrlTokens externalStage, int fileNumber) { var files = new List(); 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(); _extractService = new Mock(); _historyService = new Mock(); _recurringJobService = new Mock(); _sftpService = new Mock(); _dataQueryService = new Mock(); _claimsPrincipalAccessor = IClaimsPrincipalAccessorMock.GetClaimsPrincipalAccessor(); _logger = new Mock>(); _scheduleService = new Mock(); _config = new Mock>(); 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" }; } } }