using GB5Shared.DTO.Framework.Login; using GB5Shared.QueryExecutor; using IceImportDAL.DTO.IceImportRun; using IceImportDAL.Query.IceImportRun; namespace IceImportDAL.CustomCode.IceImportRun; public class IceImportRunDAL : IIceImportRunDAL { private readonly IQueryExecutor _qe; public IceImportRunDAL(IQueryExecutor qe) => _qe = qe; public async Task InsertRunAsync( int iceMapId, string? triggeredBy, int? sourceProfileId, long? retryOfRunId, LoginDTO login, CancellationToken ct) => await _qe.QuerySingleAsync(login, IceImportRunQB.INSERT_RUN, new { IceMapId = iceMapId, TriggeredBy = triggeredBy, SourceProfileId = sourceProfileId, RetryOfRunId = retryOfRunId, TenantId = login.ClientId, CreatedById = login.UserId }, cancellationToken: ct).ConfigureAwait(false); public async Task UpdateSourceFileArtifactAsync(long runId, long artifactId, LoginDTO login, CancellationToken ct) => await _qe.ExecuteAsync(login, IceImportRunQB.UPDATE_SOURCE_FILE_ARTIFACT, new { RunId = runId, ArtifactId = artifactId, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); public async Task UpdateStageAsync(long runId, string currentStage, LoginDTO login, CancellationToken ct) => await _qe.ExecuteAsync(login, IceImportRunQB.UPDATE_STAGE, new { RunId = runId, CurrentStage = currentStage, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); public async Task CompleteRunAsync( long runId, string status, int rowTotalCount, int successCount, int failedCount, string? errorMessage, LoginDTO login, CancellationToken ct) => await _qe.ExecuteAsync(login, IceImportRunQB.COMPLETE_RUN, new { RunId = runId, Status = status, RowTotalCount = rowTotalCount, SuccessCount = successCount, FailedCount = failedCount, ErrorMessage = errorMessage, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); public async Task CancelRunAsync(long runId, LoginDTO login, CancellationToken ct) => await _qe.ExecuteAsync(login, IceImportRunQB.CANCEL_RUN, new { RunId = runId, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); public async Task InsertLogAsync(long runId, string logLevel, string message, LoginDTO login, CancellationToken ct) => await _qe.ExecuteAsync(login, IceImportRunQB.INSERT_LOG, new { RunId = runId, LogLevel = logLevel, Message = message }, cancellationToken: ct).ConfigureAwait(false); public async Task InsertArtifactAsync(long runId, string artifactType, string storageKey, LoginDTO login, CancellationToken ct) => await _qe.QuerySingleAsync(login, IceImportRunQB.INSERT_ARTIFACT, new { RunId = runId, ArtifactType = artifactType, StorageKey = storageKey }, cancellationToken: ct).ConfigureAwait(false); public async Task InsertRowResultsAsync( long runId, IEnumerable results, LoginDTO login, CancellationToken ct) { var rows = results.Select(r => new { RunId = runId, RowNumber = r.RowNumber, EntityCode = r.EntityCode, Status = r.Status, TargetEntityId = r.TargetEntityId, PostData = r.PostData, ErrorMessage = r.ErrorMessage }).ToList(); if (rows.Count == 0) return; // >100 rows is the realistic norm for a bulk import — CLAUDE.md mandates BulkInsertAsync // over a looped ExecuteAsync for exactly this shape of write. await _qe.BulkInsertAsync(login, IceImportRunQB.INSERT_ROW_RESULT, rows).ConfigureAwait(false); } public async Task GetByIdAsync(long runId, LoginDTO login, CancellationToken ct) => await _qe.QuerySingleAsync(login, IceImportRunQB.GET_BY_ID, new { RunId = runId, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); public async Task> GetLogsAsync(long runId, LoginDTO login, CancellationToken ct) => await _qe.QueryAsync(login, IceImportRunQB.GET_LOGS_BY_RUN, new { RunId = runId }, cancellationToken: ct).ConfigureAwait(false); public async Task> GetArtifactsAsync(long runId, LoginDTO login, CancellationToken ct) => await _qe.QueryAsync(login, IceImportRunQB.GET_ARTIFACTS_BY_RUN, new { RunId = runId }, cancellationToken: ct).ConfigureAwait(false); public async Task GetLatestArtifactByTypeAsync( long runId, string artifactType, LoginDTO login, CancellationToken ct) => await _qe.QuerySingleAsync(login, IceImportRunQB.GET_LATEST_ARTIFACT_BY_TYPE, new { RunId = runId, ArtifactType = artifactType }, cancellationToken: ct).ConfigureAwait(false); public async Task> GetHistoryAsync( int iceMapId, int offset, int pageSize, LoginDTO login, CancellationToken ct) { var param = new { IceMapId = iceMapId, TenantId = login.ClientId, Offset = offset, PageSize = pageSize }; var items = await _qe.QueryAsync( login, IceImportRunQB.GET_HISTORY_PAGE, param, cancellationToken: ct).ConfigureAwait(false); var total = await _qe.ExecuteScalarAsync( login, IceImportRunQB.GET_HISTORY_COUNT, param, cancellationToken: ct).ConfigureAwait(false); return new PagedResult { Items = items.ToList(), TotalCount = total }; } public async Task> GetRowErrorsAsync( long runId, string? status, int offset, int pageSize, LoginDTO login, CancellationToken ct) { var param = new { RunId = runId, Status = status, Offset = offset, PageSize = pageSize }; var items = await _qe.QueryAsync( login, IceImportRunQB.GET_ROW_ERRORS_PAGE, param, cancellationToken: ct).ConfigureAwait(false); var total = await _qe.ExecuteScalarAsync( login, IceImportRunQB.GET_ROW_ERRORS_COUNT, param, cancellationToken: ct).ConfigureAwait(false); return new PagedResult { Items = items.ToList(), TotalCount = total }; } }