using FrameworkDAL.DTO.DataSync; using FrameworkDAL.Query.DataSync; using GB5Shared.DTO.Framework.Login; using GB5Shared.QueryExecutor; using GB5Shared.Resource.Response; using GB5Shared.ResponseStandard; using Newtonsoft.Json; namespace FrameworkDAL.CustomCode.DataSync { public class DataSyncJobDAL : IDataSyncJobDAL { private readonly IQueryExecutor _queryExecutor; public DataSyncJobDAL(IQueryExecutor queryExecutor) { _queryExecutor = queryExecutor; } public async Task GetSyncJob(int syncJobId, LoginDTO login, CancellationToken ct) { try { var result = await _queryExecutor.QuerySingleAsync( login, SyncJobQB.GET_SYNCJOB, new { SyncJobId = syncJobId }, cancellationToken: ct).ConfigureAwait(false); return JsonConvert.SerializeObject(result); } catch (Exception ex) { throw new Exception($"Failed to retrieve SyncJob: {ex.Message}", ex); } } public async Task GetSyncJobList(int offset, int pageSize, LoginDTO login, CancellationToken ct) { try { var param = new { Offset = offset, PageSize = pageSize }; var list = await _queryExecutor.QueryAsync( login, SyncJobQB.GET_SYNCJOB_LIST, param, cancellationToken: ct).ConfigureAwait(false); return JsonConvert.SerializeObject(list); } catch (Exception ex) { throw new Exception($"Failed to retrieve SyncJob list: {ex.Message}", ex); } } public async Task GetSelectListSyncJob(int firstNumber, int maxResult, LoginDTO login, CancellationToken ct) { try { var param = new { FirstNumber = firstNumber, MaxResult = maxResult }; var list = await _queryExecutor.QueryAsync( login, SyncJobQB.GET_SELECTLIST_SYNCJOB, param, cancellationToken: ct).ConfigureAwait(false); return JsonConvert.SerializeObject(list); } catch (Exception ex) { throw new Exception($"Failed to retrieve SyncJob select list: {ex.Message}", ex); } } public async Task SaveSyncJob(SyncJobDTO dto, LoginDTO login, CancellationToken ct) { try { await _queryExecutor.ExecuteAsync( login, SyncJobQB.SAVE_SYNCJOB, dto, cancellationToken: ct).ConfigureAwait(false); return dto.SyncJobId; } catch (Exception ex) { throw new Exception($"{ErrorResponse.SaveErrorMessage}: {ex.Message}", ex); } } public async Task UpdateSyncJob(SyncJobDTO dto, LoginDTO login, CancellationToken ct) { try { await _queryExecutor.ExecuteAsync( login, SyncJobQB.UPDATE_SYNCJOB, dto, cancellationToken: ct).ConfigureAwait(false); return dto.SyncJobId; } catch (Exception ex) { throw new Exception($"{ErrorResponse.UpdateErrorMessage}: {ex.Message}", ex); } } public async Task> DeleteSyncJob(int syncJobId, LoginDTO login, CancellationToken ct) { try { var rows = await _queryExecutor.ExecuteAsync( login, SyncJobQB.DELETE_SYNCJOB, new { SyncJobId = syncJobId }, cancellationToken: ct) .ConfigureAwait(false); return rows > 0 ? Result.Success(SuccessResponse.DeleteSuccessMessage) : Result.Failure("SyncJob not found or could not be deleted."); } catch (Exception ex) { throw new Exception("Error while deleting SyncJob", ex); } } public async Task> GetEnabledJobsForScheduler(LoginDTO login, CancellationToken ct) { try { return await _queryExecutor.QueryAsync( login, SyncJobQB.GET_ENABLED_JOBS_FOR_SCHEDULER, cancellationToken: ct) .ConfigureAwait(false); } catch (Exception ex) { throw new Exception($"Failed to load enabled sync jobs: {ex.Message}", ex); } } public async Task LockSyncJobForRun(int syncJobId, LoginDTO login, CancellationToken ct) { try { return await _queryExecutor.ExecuteAsync( login, SyncJobQB.LOCK_SYNCJOB_FOR_RUN, new { SyncJobId = syncJobId }, cancellationToken: ct) .ConfigureAwait(false); } catch (Exception ex) { throw new Exception($"Failed to lock SyncJob for run: {ex.Message}", ex); } } public async Task UpdateLastRun(int syncJobId, int runId, byte status, DateTime ranAt, LoginDTO login, CancellationToken ct) { try { await _queryExecutor.ExecuteAsync(login, SyncJobQB.UPDATE_LASTRUN, new { SyncJobId = syncJobId, LastRunId = runId, LastRunStatus = status, LastRunAt = ranAt, ModifiedById = login.UserId }, cancellationToken: ct).ConfigureAwait(false); } catch (Exception ex) { throw new Exception($"Failed to update last run for SyncJob {syncJobId}: {ex.Message}", ex); } } public async Task UpdateWatermarks(int syncJobId, DateTime insertedTill, DateTime updatedTill, DateTime deletedTill, LoginDTO login, CancellationToken ct) { try { await _queryExecutor.ExecuteAsync(login, SyncJobQB.UPDATE_WATERMARKS, new { SyncJobId = syncJobId, LastInsertedTill = insertedTill, LastUpdatedTill = updatedTill, LastDeleteTill = deletedTill, ModifiedById = login.UserId }, cancellationToken: ct).ConfigureAwait(false); } catch (Exception ex) { throw new Exception($"Failed to update watermarks for SyncJob {syncJobId}: {ex.Message}", ex); } } } }