Я создал систему очередей за webhook. Запросы поступают из устройств IoT и их данных, которые были отправлены в отправку по почте с использованием Hangefire (используется для управления повторными). Только 1 работа в сфере Hangfire Per Webhookid будет обрабатываться в любое время. Я настроил тестовый клиент для получения Webhooks, и я возвращаю статус HTTP 400 для тестирования повторений. После того, как работа примерно на 10 секунд преуспевает - хотя на моей работе перечислены httprequestexception (badrequest) на панели панели. < /P>
Любые идеи, пожалуйста?public class WebhookQueueService
{
private readonly ConcurrentDictionary _queues = new();
private readonly ConcurrentDictionary _jobPending = new();
private readonly IBackgroundJobClient _backgroundJobClient;
public WebhookQueueService(IBackgroundJobClient backgroundJobClient)
{
_backgroundJobClient = backgroundJobClient;
}
public void EnqueueEvents(string webhookId, List events)
{
var queue = _queues.GetOrAdd(webhookId, _ => new ConcurrentQueue());
foreach (var ev in events)
{
queue.Enqueue(ev);
}
// If no job is pending for this webhook, queue one
if (_jobPending.TryAdd(webhookId, true))
{
_backgroundJobClient.Enqueue(p =>
p.ProcessWebhookQueueAsync(webhookId));
}
}
public List DequeueAll(string webhookId)
{
var result = new List();
if (_queues.TryGetValue(webhookId, out var queue))
{
while (queue.TryDequeue(out var ev))
{
result.Add(ev);
}
}
return result;
}
public void MarkJobComplete(string webhookId)
{
_jobPending.TryRemove(webhookId, out _);
}
}
< /code>
webhookqueueprocessor: < /p>
[AutomaticRetry(Attempts = 5)]
public async Task ProcessWebhookQueueAsync(string webhookId)
{
try
{
var events = _queueService.DequeueAll(webhookId);
if (events.Count == 0) return;
if (!int.TryParse(webhookId, out var id)) return;
var webhook = await _dbContext.Webhook
.Where(w => w.Id == id && w.IsActive)
.FirstOrDefaultAsync();
if (webhook == null) return;
var payload = new WebhookPayload { Events = events };
await _webhookService.SendLogToWebhookAsync(webhook.EndpointUrl, payload, webhook.SecretToken);
// Only mark complete after successful send
_queueService.MarkJobComplete(webhookId);
}
catch (Exception ex)
{
Console.WriteLine("Error processing webhook queue for {0}: {1}", webhookId, ex.Message);
throw; // triggers Hangfire retry
}
}
< /code>
webhookservice: < /p>
public async Task SendLogToWebhookAsync(string endpointUrl, WebhookPayload payload, string? secret = null)
{
try
{
var client = _httpClientFactory.CreateClient("WebhookClient");
var json = JsonSerializer.Serialize(payload);
var request = new HttpRequestMessage(HttpMethod.Post, endpointUrl)
{
Content = new StringContent(json, Encoding.UTF8, "application/json")
};
if (!string.IsNullOrEmpty(secret))
{
request.Headers.Add("X-Webhook-Secret", secret);
}
var response = await client.SendAsync(request);
if (!response.IsSuccessStatusCode)
{
Console.WriteLine("THROWING due to status {0}", response.StatusCode);
throw new HttpRequestException($"Non-success status code: {response.StatusCode}");
}
response.EnsureSuccessStatusCode(); // triggers Hangfire retry if not 2xx
}
catch (HttpRequestException ex)
{
Console.WriteLine("Webhook POST failed to {0}. Reason: {1}", endpointUrl, ex.Message);
throw; // let Hangfire retry
}
catch (Exception ex)
{
Console.WriteLine("Unexpected exception while sending webhook to {0}. Reason: {1}", endpointUrl, ex.Message);
throw; // let Hangfire retry
}
}
< /code>
Программа: < /p>
builder.Services.AddSingleton();
builder.Services.AddScoped();
builder.Services.AddScoped();
Подробнее здесь: https://stackoverflow.com/questions/796 ... texception