/
rezvich
/
adapter_v
Обзор
Документация
Войти
/
rezvich
/
adapter_v
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
Repositories/OutboxRepository.cs
166 строк
6 KB
ItsEnd
HttpClientFactory для SOAP
21 дек 2025, 18:15
21 дек 2025, 18:15
a4a1c0f
Код
Авторство
О чём код?
using Microsoft.Data.SqlClient; using Microsoft.Extensions.Configuration; namespace adapter_v.Repositories { public interface IOutboxRepository { Task<IReadOnlyList<SoapOutboxItem>> GetBatchForSendAsync( int batchSize, int leaseMinutes, int maxTries, string lockedBy, CancellationToken ct = default); Task MarkSentAsync(long id, int httpStatus, CancellationToken ct = default); Task MarkRetryAsync(long id, int? httpStatus, string error, int retryDelaySeconds, CancellationToken ct = default); Task MarkErrorFinalAsync(long id, int? httpStatus, string error, CancellationToken ct = default); } public sealed class OutboxRepository : IOutboxRepository { private readonly string _cs; public OutboxRepository(IConfiguration cfg) { _cs = cfg.GetConnectionString("Db") ?? throw new InvalidOperationException("ConnectionStrings:Db not found"); } public async Task<IReadOnlyList<SoapOutboxItem>> GetBatchForSendAsync( int batchSize, int leaseMinutes, int maxTries, string lockedBy, CancellationToken ct = default) { const string sql = @" ;WITH cte AS ( SELECT TOP (@batchSize) o.id FROM dbo.SoapOutbox o WITH (READPAST, UPDLOCK, ROWLOCK) WHERE (o.status = 'NEW' OR (o.status = 'IN_PROGRESS' AND o.lease_until IS NOT NULL AND o.lease_until < SYSUTCDATETIME())) AND (o.next_attempt_at IS NULL OR o.next_attempt_at <= SYSUTCDATETIME()) AND o.try_count < @maxTries ORDER BY o.id ) UPDATE o SET o.status = 'IN_PROGRESS', o.try_count = o.try_count + 1, o.updated_at = SYSUTCDATETIME(), o.last_attempt_at = SYSUTCDATETIME(), o.external_request_id = COALESCE(o.external_request_id, NEWID()), o.locked_by = @lockedBy, o.lease_until = DATEADD(minute, @leaseMinutes, SYSUTCDATETIME()), o.next_attempt_at = NULL OUTPUT inserted.id, inserted.external_request_id, inserted.url, inserted.soap_xml, inserted.try_count FROM dbo.SoapOutbox o JOIN cte ON cte.id = o.id; "; var result = new List<SoapOutboxItem>(batchSize); await using var con = new SqlConnection(_cs); await con.OpenAsync(ct); await using var cmd = new SqlCommand(sql, con); cmd.Parameters.Add(new SqlParameter("@batchSize", System.Data.SqlDbType.Int) { Value = batchSize }); cmd.Parameters.Add(new SqlParameter("@leaseMinutes", System.Data.SqlDbType.Int) { Value = leaseMinutes }); cmd.Parameters.Add(new SqlParameter("@maxTries", System.Data.SqlDbType.Int) { Value = maxTries }); cmd.Parameters.Add(new SqlParameter("@lockedBy", System.Data.SqlDbType.NVarChar, 128) { Value = lockedBy }); await using var rd = await cmd.ExecuteReaderAsync(ct); while (await rd.ReadAsync(ct)) { result.Add(new SoapOutboxItem( rd.GetInt64(0), rd.IsDBNull(1) ? null : rd.GetGuid(1), rd.GetString(2), rd.GetString(3), rd.GetInt32(4) )); } return result; } public async Task MarkSentAsync(long id, int httpStatus, CancellationToken ct = default) { const string sql = @" UPDATE dbo.SoapOutbox SET status = 'SENT', updated_at = SYSUTCDATETIME(), last_http_status = @http_status, last_error = NULL, lease_until = NULL, locked_by = NULL, next_attempt_at = NULL WHERE id = @id; "; await ExecAsync(sql, ct, new SqlParameter("@id", System.Data.SqlDbType.BigInt) { Value = id }, new SqlParameter("@http_status", System.Data.SqlDbType.Int) { Value = httpStatus }); } public async Task MarkRetryAsync(long id, int? httpStatus, string error, int retryDelaySeconds, CancellationToken ct = default) { const string sql = @" UPDATE dbo.SoapOutbox SET status = 'NEW', updated_at = SYSUTCDATETIME(), last_http_status = @http_status, last_error = @err, lease_until = NULL, locked_by = NULL, next_attempt_at = DATEADD(second, @delay_sec, SYSUTCDATETIME()) WHERE id = @id; "; await ExecAsync(sql, ct, new SqlParameter("@id", System.Data.SqlDbType.BigInt) { Value = id }, new SqlParameter("@http_status", System.Data.SqlDbType.Int) { Value = (object?)httpStatus ?? DBNull.Value }, new SqlParameter("@err", System.Data.SqlDbType.NVarChar, -1) { Value = (object?)error ?? DBNull.Value }, new SqlParameter("@delay_sec", System.Data.SqlDbType.Int) { Value = retryDelaySeconds }); } public async Task MarkErrorFinalAsync(long id, int? httpStatus, string error, CancellationToken ct = default) { const string sql = @" UPDATE dbo.SoapOutbox SET status = 'ERROR_FINAL', updated_at = SYSUTCDATETIME(), last_http_status = @http_status, last_error = @err, lease_until = NULL, locked_by = NULL WHERE id = @id; "; await ExecAsync(sql, ct, new SqlParameter("@id", System.Data.SqlDbType.BigInt) { Value = id }, new SqlParameter("@http_status", System.Data.SqlDbType.Int) { Value = (object?)httpStatus ?? DBNull.Value }, new SqlParameter("@err", System.Data.SqlDbType.NVarChar, -1) { Value = (object?)error ?? DBNull.Value }); } private async Task ExecAsync(string sql, CancellationToken ct, params SqlParameter[] parameters) { await using var con = new SqlConnection(_cs); await con.OpenAsync(ct); await using var cmd = new SqlCommand(sql, con); if (parameters is { Length: > 0 }) cmd.Parameters.AddRange(parameters); await cmd.ExecuteNonQueryAsync(ct); } } }