using GB5Shared.DTO.Framework.Login; using GB5Shared.QueryExecutor; using JobEngineDAL.DTOs; using JobEngineDAL.Interfaces; using JobEngineDAL.Query; namespace JobEngineDAL.Implementations { public class JobExecutionDAL : IJobExecutionDAL { private readonly IQueryExecutor _qe; public JobExecutionDAL(IQueryExecutor queryExecutor) => _qe = queryExecutor; public async Task InsertExecutionAsync(JobExecutionDTO dto, LoginDTO login, CancellationToken ct) => await _qe.ExecuteScalarAsync(login, JobExecutionQB.INSERT_EXECUTION, new { dto.JobId, TenantId = login.ClientId, dto.JobName, dto.TriggerType, dto.ScheduledFor, dto.RetryAttempt, dto.CorrelationId, HostInstance = Environment.MachineName }, cancellationToken: ct).ConfigureAwait(false); public async Task CompleteExecutionAsync(long executionId, string status, string? resultJson, string? errorDetails, LoginDTO login, CancellationToken ct) => await _qe.ExecuteAsync(login, JobExecutionQB.COMPLETE_EXECUTION, new { ExecutionId = executionId, Status = status, ResultJson = resultJson, ErrorDetails = errorDetails, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); public async Task<(IEnumerable Items, int TotalCount)> GetHistoryPagedAsync( int tenantId, int? jobId, string? status, int page, int pageSize, LoginDTO login, CancellationToken ct) { var whereClause = ""; if (jobId.HasValue) whereClause += " AND JOBID = @JobId"; if (!string.IsNullOrWhiteSpace(status)) whereClause += " AND STATUS = @Status"; var sql = JobExecutionQB.GET_HISTORY_PAGED.Replace("{WHERE_CLAUSE}", whereClause); var countSql = JobExecutionQB.GET_HISTORY_COUNT.Replace("{WHERE_CLAUSE}", whereClause); var param = new { TenantId = tenantId, JobId = jobId, Status = status, Offset = (page - 1) * pageSize, PageSize = pageSize }; var items = await _qe.QueryAsync(login, sql, param, cancellationToken: ct).ConfigureAwait(false); var total = await _qe.ExecuteScalarAsync(login, countSql, new { TenantId = tenantId, JobId = jobId, Status = status }, cancellationToken: ct).ConfigureAwait(false); return (items, total); } public async Task GetExecutionDetailAsync(long executionId, LoginDTO login, CancellationToken ct) => await _qe.QuerySingleAsync(login, JobExecutionQB.GET_EXECUTION_DETAIL, new { ExecutionId = executionId, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); public async Task> GetRunningJobsAsync(LoginDTO login, CancellationToken ct) => await _qe.QueryAsync(login, JobExecutionQB.GET_RUNNING_JOBS, new { TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); public async Task<(int SuccessCount, int FailedCount, int TotalCount)> GetStats7DaysAsync(LoginDTO login, CancellationToken ct) { var rows = await _qe.QueryAsync(login, JobExecutionQB.GET_STATS_7DAYS, new { TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); var row = rows.FirstOrDefault(); if (row == null) return (0, 0, 0); return (Convert.ToInt32(row.SuccessCount), Convert.ToInt32(row.FailedCount), Convert.ToInt32(row.TotalCount)); } public async Task> GetAvgDurationsAsync(LoginDTO login, CancellationToken ct) => await _qe.QueryAsync(login, JobExecutionQB.GET_AVG_DURATIONS, new { TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); } }