using System; using System.Collections.Generic; using System.IO; using System.Text; using System.Text.Json; using System.Threading; using System.Threading.Tasks; using DMSDAL.CustomCode.BulkIngestion; using DMSDAL.DTO.BulkIngestion; using GB5Shared.Attachment; using GB5Shared.DTO.ECM; using GB5Shared.DTO.Framework.Login; using Microsoft.Extensions.Logging; namespace DMSBLL.BulkIngestion { public class BulkIngestionBLL : IBulkIngestionBLL { private readonly IBulkIngestionDAL _dal; private readonly IAttachmentUploadService _upload; private readonly ILogger _log; // Batch size for DB job counter updates private const int CounterFlushInterval = 10; public BulkIngestionBLL( IBulkIngestionDAL dal, IAttachmentUploadService upload, ILogger log) { _dal = dal; _upload = upload; _log = log; } public async Task StartIngestionAsync( string sourceFolder, string? manifestFilePath, int documentSetDetailId, LoginDTO login, CancellationToken ct = default) { if (!Directory.Exists(sourceFolder)) throw new DirectoryNotFoundException($"Source folder not found: {sourceFolder}"); // ── Mode detection ──────────────────────────────────── var effectiveManifest = manifestFilePath ?? Path.Combine(sourceFolder, "manifest.csv"); var useManifest = File.Exists(effectiveManifest); var mode = useManifest ? (byte)0 : (byte)1; // ── Create job row ──────────────────────────────────── var job = new BulkIngestionJobDTO { SourceFolder = sourceFolder, ManifestFilePath = useManifest ? effectiveManifest : null, DocumentSetDetailId = documentSetDetailId, Mode = mode, Status = 0, CreatedByUserId = login.UserId }; var jobId = await _dal.InsertJobAsync(job, login, ct).ConfigureAwait(false); // ── Build manifest rows ─────────────────────────────── var rows = useManifest ? ParseManifestCsv(effectiveManifest) : await ScanFolderConventionAsync(sourceFolder, login, ct); var total = 0; var processed = 0; var failed = 0; foreach (var row in rows) total++; await _dal.UpdateJobStatusAsync(jobId, 1, total, 0, 0, login, ct) .ConfigureAwait(false); // ── Ingest each file ────────────────────────────────── int batchCounter = 0; foreach (var row in useManifest ? ParseManifestCsv(effectiveManifest) : await ScanFolderConventionAsync(sourceFolder, login, ct)) { ct.ThrowIfCancellationRequested(); byte logStatus = 1; string? errorMsg = null; int? attachmentId = null; try { if (!File.Exists(row.FilePath)) throw new FileNotFoundException($"File not found: {row.FilePath}"); var contextDict = string.IsNullOrWhiteSpace(row.ContextJson) ? new Dictionary() : JsonSerializer.Deserialize>(row.ContextJson) ?? new Dictionary(); var request = new AttachmentUploadRequest { DocumentSetDetailId = row.DocumentSetDetailId > 0 ? row.DocumentSetDetailId : documentSetDetailId, ObjectTypeId = row.ObjectTypeId, ObjectId = row.ObjectId, BizTransactionTypeId = -1, RowGuid = Guid.NewGuid(), Tags = row.Tags ?? string.Empty, ContextDictionary = contextDict }; var mimeType = MimeFromExtension(Path.GetExtension(row.FilePath)); await using var fs = new FileStream( row.FilePath, FileMode.Open, FileAccess.Read, FileShare.Read, 81_920, useAsync: true); var result = await _upload.UploadAsync( request, fs, Path.GetFileName(row.FilePath), mimeType, login, ct) .ConfigureAwait(false); attachmentId = result.AttachmentId; logStatus = 1; // Done processed++; } catch (Exception ex) { _log.LogError(ex, "Bulk ingestion file failed. Job={Job} File={File}", jobId, row.FilePath); errorMsg = ex.Message; logStatus = 2; // Failed failed++; } await _dal.InsertLogAsync( jobId, row.FilePath, row.ObjectTypeId, row.ObjectId, row.DocumentSetDetailId, row.DisplayName, attachmentId, logStatus, errorMsg, login, ct).ConfigureAwait(false); if (++batchCounter % CounterFlushInterval == 0) await _dal.UpdateJobStatusAsync( jobId, 1, total, processed, failed, login, ct).ConfigureAwait(false); } var finalStatus = (byte)(failed == total && total > 0 ? 3 : 2); await _dal.UpdateJobStatusAsync( jobId, finalStatus, total, processed, failed, login, ct).ConfigureAwait(false); return new BulkIngestionStatusDTO { JobId = jobId, Status = finalStatus, TotalFiles = total, ProcessedFiles = processed, FailedFiles = failed }; } public async Task GetIngestionStatusAsync( int jobId, LoginDTO login, CancellationToken ct = default) => await _dal.GetJobStatusAsync(jobId, login, ct).ConfigureAwait(false); // ── Manifest CSV parsing ───────────────────────────────────────── // Columns: FilePath,ObjectTypeId,ObjectId,DocumentSetDetailId,DisplayName,Tags,ContextJson // First row is header — skipped. private static IEnumerable ParseManifestCsv(string path) { using var reader = new StreamReader(path, Encoding.UTF8); string? line; bool headerSkipped = false; while ((line = reader.ReadLine()) is not null) { if (!headerSkipped) { headerSkipped = true; continue; } if (string.IsNullOrWhiteSpace(line)) continue; var cols = line.Split(',', 7); if (cols.Length < 3) continue; yield return new BulkIngestionManifestRow { FilePath = cols[0].Trim('"', ' '), ObjectTypeId = int.TryParse(cols[1].Trim(), out var ot) ? ot : 0, ObjectId = int.TryParse(cols[2].Trim(), out var oi) ? oi : 0, DocumentSetDetailId = cols.Length > 3 && int.TryParse(cols[3].Trim(), out var ds) ? ds : -1, DisplayName = cols.Length > 4 ? cols[4].Trim('"', ' ') : null, Tags = cols.Length > 5 ? cols[5].Trim('"', ' ') : null, ContextJson = cols.Length > 6 ? cols[6].Trim('"', ' ') : null }; } } // ── Folder convention scanning ─────────────────────────────────── // Expected layout: {sourceFolder}/{ObjectTypeCode}/{ObjectId}/file.ext private async Task> ScanFolderConventionAsync( string sourceFolder, LoginDTO login, CancellationToken ct) { var rows = new List(); foreach (var typeDir in Directory.EnumerateDirectories(sourceFolder)) { var typeCode = Path.GetFileName(typeDir); var typeId = await _dal.GetObjectTypeIdByCodeAsync(typeCode, login, ct) .ConfigureAwait(false); if (typeId is null) { _log.LogWarning("Folder convention: unknown EntityCode '{Code}' — skipping folder.", typeCode); continue; } foreach (var objDir in Directory.EnumerateDirectories(typeDir)) { if (!int.TryParse(Path.GetFileName(objDir), out var objectId)) continue; foreach (var file in Directory.EnumerateFiles(objDir)) { rows.Add(new BulkIngestionManifestRow { FilePath = file, ObjectTypeId = typeId.Value, ObjectId = objectId, DisplayName = Path.GetFileNameWithoutExtension(file) }); } } } return rows; } private static string MimeFromExtension(string ext) => ext.ToLowerInvariant() switch { ".pdf" => "application/pdf", ".png" => "image/png", ".jpg" => "image/jpeg", ".jpeg" => "image/jpeg", ".xlsx" => "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet", ".docx" => "application/vnd.openxmlformats-officedocument.wordprocessingml.document", ".txt" => "text/plain", ".csv" => "text/csv", _ => "application/octet-stream" }; } }