Skip to content

Commit 3db02ac

Browse files
authored
LNK-4430: Data Acquisition - Handle 429 Status (#1378)
* Checkin * Fix tests * Don't Count 429 as a Failure * ParseRetryAfter refactor and safety * refactor * refactor * final refactor * Add More Tests * ExecutionDate filtering in query * 429 Set ExecutionDate for ALL of the facilities Pending, Ready, and Failed logs with an execution time less then or equal to the RetryAfter derived time. * changes after convo
1 parent 6c1cacb commit 3db02ac

17 files changed

Lines changed: 557 additions & 61 deletions

File tree

DotNet/DataAcquisition.Domain/Application/Managers/DataAcquisitionLogManager.cs

Lines changed: 32 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,13 @@
1-
using System.Diagnostics;
21
using DataAcquisition.Domain.Application.Models;
32
using LantanaGroup.Link.DataAcquisition.Domain.Application.Models.Api.QueryLog;
43
using LantanaGroup.Link.DataAcquisition.Domain.Application.Models.Exceptions;
4+
using LantanaGroup.Link.DataAcquisition.Domain.Application.Queries;
55
using LantanaGroup.Link.DataAcquisition.Domain.Infrastructure;
66
using LantanaGroup.Link.DataAcquisition.Domain.Infrastructure.Entities;
77
using LantanaGroup.Link.DataAcquisition.Domain.Infrastructure.Models.Enums;
88
using LantanaGroup.Link.DataAcquisition.Domain.Models;
9-
using LantanaGroup.Link.Shared.Application.Models.Telemetry;
109
using Microsoft.Extensions.Logging;
10+
using System.Diagnostics;
1111

1212
namespace LantanaGroup.Link.DataAcquisition.Domain.Application.Managers;
1313

@@ -17,17 +17,20 @@ public interface IDataAcquisitionLogManager
1717
Task<DataAcquisitionLogModel?> UpdateAsync(UpdateDataAcquisitionLogModel updateLog, CancellationToken cancellationToken = default);
1818
Task DeleteAsync(long id, CancellationToken cancellationToken = default);
1919
Task UpdateTailFlagForFacilityCorrelationIdReportTrackingId(List<long> logIds, string facilityId, string correlationId, string reportTrackingId, CancellationToken cancellationToken = default);
20+
Task ThrottleFacilityAcquisitions(string facilityId, DateTime executionDate, CancellationToken cancellationToken = default);
2021
}
2122

2223
public class DataAcquisitionLogManager : IDataAcquisitionLogManager
2324
{
2425
public readonly ILogger<DataAcquisitionLogManager> _logger;
2526
public readonly IDatabase _database;
27+
private readonly IDataAcquisitionLogQueries _logQueries;
2628

27-
public DataAcquisitionLogManager(ILogger<DataAcquisitionLogManager> logger, IDatabase database)
29+
public DataAcquisitionLogManager(ILogger<DataAcquisitionLogManager> logger, IDatabase database, IDataAcquisitionLogQueries logQueries)
2830
{
2931
_logger = logger ?? throw new ArgumentNullException(nameof(logger));
3032
_database = database ?? throw new ArgumentNullException(nameof(database));
33+
_logQueries = logQueries;
3134
}
3235

3336
public async Task<DataAcquisitionLogModel> CreateAsync(CreateDataAcquisitionLogModel model, CancellationToken cancellationToken = default)
@@ -202,4 +205,30 @@ public async Task UpdateTailFlagForFacilityCorrelationIdReportTrackingId(List<lo
202205
await _database.DataAcquisitionLogRepository.SaveChangesAsync();
203206
}
204207
}
208+
209+
public async Task ThrottleFacilityAcquisitions(string facilityId, DateTime executionDate, CancellationToken cancellationToken = default)
210+
{
211+
long? lastId = null;
212+
var batchSize = 1000;
213+
while (true)
214+
{
215+
// Get All Active Logs For Batch
216+
var toThrottle = await _logQueries.GetNextEligibleBatchForFacility(facilityId, lastId, batchSize, [RequestStatus.Failed, RequestStatus.Ready, RequestStatus.Pending], executionDate, cancellationToken);
217+
218+
if(toThrottle.Count == 0)
219+
{
220+
break;
221+
}
222+
223+
//Update their next processing time
224+
foreach (var log in toThrottle)
225+
{
226+
log.ExecutionDate = executionDate;
227+
}
228+
229+
await _database.DataAcquisitionLogRepository.SaveChangesAsync();
230+
231+
lastId = toThrottle.Max(l => l.Id);
232+
}
233+
}
205234
}
Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
namespace LantanaGroup.Link.DataAcquisition.Domain.Application.Models.Exceptions
2+
{
3+
public class ProcessingDelayException : Exception
4+
{
5+
public ProcessingDelayException(string message) : base(message) { }
6+
}
7+
}
Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,14 @@
1+
using System;
2+
using Hl7.Fhir.Rest;
3+
4+
namespace LantanaGroup.Link.DataAcquisition.Domain.Application.Models.Exceptions
5+
{
6+
public class TooManyRequestsException : FhirOperationException
7+
{
8+
public TimeSpan RetryAfter { get; }
9+
public TooManyRequestsException(string message, TimeSpan retryAfter) : base(message, System.Net.HttpStatusCode.TooManyRequests)
10+
{
11+
RetryAfter = retryAfter;
12+
}
13+
}
14+
}

DotNet/DataAcquisition.Domain/Application/Queries/DataAcquisitionLogQueries.cs

Lines changed: 12 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -67,7 +67,7 @@ public interface IDataAcquisitionLogQueries
6767

6868
Task<List<string>> GetFacilitiesWithPendingAndRetryableFailedRequests(CancellationToken cancellationToken = default);
6969

70-
Task<List<DataAcquisitionLogModel>> GetNextEligibleBatchForFacility(string facilityId, long? lastId, int batchSize, CancellationToken cancellationToken = default);
70+
Task<List<DataAcquisitionLogModel>> GetNextEligibleBatchForFacility(string facilityId, long? lastId, int batchSize, List<RequestStatus> statuses, DateTime? designagtedExecutionTime = null, CancellationToken cancellationToken = default);
7171

7272
Task<List<string>> GetResourceIdsForReportPatient(string correlationId, string facilityId, string resourceType, CancellationToken cancellationToken = default);
7373
}
@@ -566,13 +566,16 @@ public async Task<List<string>> GetFacilitiesWithPendingAndRetryableFailedReques
566566
.ToListAsync(cancellationToken);
567567
}
568568

569-
public async Task<List<DataAcquisitionLogModel>> GetNextEligibleBatchForFacility(string facilityId, long? lastId, int batchSize, CancellationToken cancellationToken = default)
569+
public async Task<List<DataAcquisitionLogModel>> GetNextEligibleBatchForFacility(string facilityId, long? lastId, int batchSize, List<RequestStatus> statuses, DateTime? designagtedExecutionTime = null, CancellationToken cancellationToken = default)
570570
{
571+
designagtedExecutionTime ??= DateTime.UtcNow;
572+
571573
var query = from log in _dbContext.DataAcquisitionLogs
572-
orderby log.Id
573574
where log.FacilityId == facilityId
574575
&& (lastId == null || log.Id > lastId)
575-
&& (log.Status == RequestStatus.Pending || log.Status == RequestStatus.Failed)
576+
&& (log.ExecutionDate == null || log.ExecutionDate <= designagtedExecutionTime)
577+
&& (log.Status == null || statuses.Contains(log.Status.Value))
578+
orderby log.Id
576579
select new DataAcquisitionLogModel
577580
{
578581
Id = log.Id,
@@ -622,4 +625,9 @@ private Expression<Func<T, object>> SetSortBy<T>(string? sortBy)
622625
var converted = Expression.Convert(property, typeof(object));
623626
return Expression.Lambda<Func<T, object>>(converted, parameter);
624627
}
628+
629+
public Task<List<DataAcquisitionLogModel>> GetNextBatchForFacility(string facilityId, long? lastId, int batchSize, CancellationToken cancellationToken = default)
630+
{
631+
throw new NotImplementedException();
632+
}
625633
}
Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,51 @@
1+
using System.Net.Http.Headers;
2+
3+
namespace LantanaGroup.Link.DataAcquisition.Domain.Application.Services.FhirApi.Commands;
4+
5+
internal class HeaderCapturingHandler : DelegatingHandler
6+
{
7+
public HttpResponseHeaders? LastResponseHeaders { get; private set; }
8+
9+
protected override async Task<HttpResponseMessage> SendAsync(HttpRequestMessage request, CancellationToken cancellationToken)
10+
{
11+
var response = await base.SendAsync(request, cancellationToken);
12+
LastResponseHeaders = response.Headers;
13+
return response;
14+
}
15+
}
16+
17+
public static class FhirCommandUtils
18+
{
19+
public static TimeSpan ParseRetryAfter(HttpResponseHeaders? headers, TimeSpan defaultDelay = default)
20+
{
21+
defaultDelay = defaultDelay == default ? TimeSpan.FromSeconds(60) : defaultDelay;
22+
23+
if (headers == null || headers.RetryAfter == null)
24+
{
25+
return defaultDelay;
26+
}
27+
28+
var retryValue = headers.RetryAfter;
29+
TimeSpan delay;
30+
31+
if (retryValue.Delta.HasValue)
32+
{
33+
delay = retryValue.Delta.Value;
34+
}
35+
else if (retryValue.Date.HasValue)
36+
{
37+
delay = retryValue.Date.Value - DateTimeOffset.UtcNow;
38+
}
39+
else
40+
{
41+
return defaultDelay;
42+
}
43+
44+
if (delay <= TimeSpan.Zero)
45+
{
46+
return defaultDelay;
47+
}
48+
49+
return delay;
50+
}
51+
}

DotNet/DataAcquisition.Domain/Application/Services/FhirApi/Commands/ReadFhirCommand.cs

Lines changed: 23 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,15 @@
11
using Hl7.Fhir.Model;
22
using Hl7.Fhir.Rest;
3+
using LantanaGroup.Link.DataAcquisition.Domain.Application.Factories.Auth;
34
using LantanaGroup.Link.DataAcquisition.Domain.Application.Interfaces;
45
using LantanaGroup.Link.DataAcquisition.Domain.Application.Models;
6+
using LantanaGroup.Link.DataAcquisition.Domain.Application.Models.Exceptions;
57
using LantanaGroup.Link.Shared.Application.Models.Configs;
68
using Medallion.Threading;
79
using Microsoft.Extensions.Logging;
810
using Microsoft.Extensions.Options;
11+
using System.Net;
912
using System.Net.Http.Headers;
10-
using LantanaGroup.Link.DataAcquisition.Domain.Application.Factories.Auth;
1113

1214
namespace LantanaGroup.Link.DataAcquisition.Domain.Application.Services.FhirApi.Commands;
1315

@@ -18,12 +20,13 @@ public record ReadFhirCommandRequest(
1820
string baseUrl,
1921
FhirQueryConfigurationModel fhirQueryConfiguration);
2022

21-
public interface IReadFhirCommand
22-
{
23+
public interface IReadFhirCommand
24+
{
2325
Task<DomainResource> ExecuteAsync(
2426
ReadFhirCommandRequest request,
2527
CancellationToken cancellationToken = default);
2628
}
29+
2730
public class ReadFhirCommand : IReadFhirCommand
2831
{
2932
private readonly ILogger<ReadFhirCommand> _logger;
@@ -48,21 +51,23 @@ public ReadFhirCommand(
4851

4952
public async Task<DomainResource> ExecuteAsync(ReadFhirCommandRequest request, CancellationToken cancellationToken = default)
5053
{
51-
5254
if (string.IsNullOrWhiteSpace(request.resourceId))
5355
throw new ArgumentNullException(nameof(request.resourceId), "Resource ID cannot be null or empty.");
5456

5557
if (string.IsNullOrWhiteSpace(request.baseUrl))
5658
throw new ArgumentNullException(nameof(request.baseUrl), "FhirClient Endpoint cannot be null.");
5759

58-
5960
if (request.fhirQueryConfiguration == null)
6061
throw new ArgumentNullException(nameof(request.fhirQueryConfiguration), "FhirQueryConfiguration cannot be null.");
6162

62-
6363
using (_distributedSemaphoreProvider.AcquireSemaphore(request.facilityId, request.fhirQueryConfiguration.MaxConcurrentRequests.GetValueOrDefault(1), _distributedLockSettings.Expiration, cancellationToken))
6464
{
65-
var fhirClient = new FhirClient(request.baseUrl.Trim('/'), _httpClient, new FhirClientSettings
65+
// Create a new handler chain using a DelegatingHandler around a base HttpClientHandler
66+
var innerHandler = new HttpClientHandler();
67+
var headerCapturingHandler = new HeaderCapturingHandler { InnerHandler = innerHandler };
68+
var httpClientWithHandler = new HttpClient(headerCapturingHandler);
69+
70+
var fhirClient = new FhirClient(request.baseUrl.Trim('/'), httpClientWithHandler, new FhirClientSettings
6671
{
6772
PreferredFormat = ResourceFormat.Json
6873
});
@@ -73,15 +78,23 @@ public async Task<DomainResource> ExecuteAsync(ReadFhirCommandRequest request, C
7378
fhirClient.RequestHeaders.Authorization = (AuthenticationHeaderValue)authBuilderResults.authHeader;
7479
}
7580

76-
7781
string location = request.resourceType switch
7882
{
7983
ResourceType.List => $"List/{request.resourceId}",
8084
//ResourceType.Patient => TEMPORARYPatientIdPart(id),
8185
_ => $"{request.resourceType}/{request.resourceId}"
8286
};
8387

84-
var readResource = await fhirClient.ReadAsync<DomainResource>(location);
88+
DomainResource readResource;
89+
try
90+
{
91+
readResource = await fhirClient.ReadAsync<DomainResource>(location);
92+
}
93+
catch (FhirOperationException ex) when (ex.Status == HttpStatusCode.TooManyRequests)
94+
{
95+
var retryAfter = FhirCommandUtils.ParseRetryAfter(headerCapturingHandler.LastResponseHeaders);
96+
throw new TooManyRequestsException($"Too many requests for {location}", retryAfter);
97+
}
8598

8699
if (readResource == null)
87100
{
@@ -91,4 +104,4 @@ public async Task<DomainResource> ExecuteAsync(ReadFhirCommandRequest request, C
91104
return readResource;
92105
}
93106
}
94-
}
107+
}

DotNet/DataAcquisition.Domain/Application/Services/FhirApi/Commands/SearchFhirCommand.cs

Lines changed: 38 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1,16 +1,18 @@
11
using Hl7.Fhir.Model;
22
using Hl7.Fhir.Rest;
3+
using LantanaGroup.Link.DataAcquisition.Domain.Application.Factories.Auth;
34
using LantanaGroup.Link.DataAcquisition.Domain.Application.Interfaces;
45
using LantanaGroup.Link.DataAcquisition.Domain.Application.Models;
6+
using LantanaGroup.Link.DataAcquisition.Domain.Application.Models.Exceptions;
57
using LantanaGroup.Link.DataAcquisition.Domain.Infrastructure.Models.Enums;
68
using LantanaGroup.Link.Shared.Application.Models.Configs;
79
using LantanaGroup.Link.Shared.Application.Models.Telemetry;
810
using LantanaGroup.Link.Shared.Application.Services.Security;
911
using Medallion.Threading;
1012
using Microsoft.Extensions.Logging;
1113
using Microsoft.Extensions.Options;
14+
using System.Net;
1215
using System.Net.Http.Headers;
13-
using LantanaGroup.Link.DataAcquisition.Domain.Application.Factories.Auth;
1416
using ResourceType = Hl7.Fhir.Model.ResourceType;
1517

1618
namespace LantanaGroup.Link.DataAcquisition.Domain.Application.Services.FhirApi.Commands;
@@ -22,7 +24,7 @@ public record SearchFhirCommandRequest(
2224
string? facilityId,
2325
string? patientId,
2426
string? correlationId,
25-
QueryPhase? queryPhase,
27+
QueryPhase? queryPhase,
2628
FhirQueryType queryType
2729
);
2830

@@ -65,7 +67,7 @@ public async IAsyncEnumerable<Bundle> ExecuteAsync(SearchFhirCommandRequest requ
6567
new KeyValuePair<string, object?>(DiagnosticNames.Resource, request.resourceType)
6668
]);
6769

68-
if(request == null || string.IsNullOrWhiteSpace(request.facilityId) || string.IsNullOrWhiteSpace(request.queryConfig.FhirServerBaseUrl))
70+
if (request == null || string.IsNullOrWhiteSpace(request.facilityId) || string.IsNullOrWhiteSpace(request.queryConfig.FhirServerBaseUrl))
6971
{
7072
_logger.LogError("Invalid request parameters. FacilityId: {FacilityId}; FhirServerBaseUrl: {FhirServerBaseUrl}", request?.facilityId?.Sanitize(), request?.queryConfig.FhirServerBaseUrl.Sanitize());
7173
yield break;
@@ -74,8 +76,12 @@ public async IAsyncEnumerable<Bundle> ExecuteAsync(SearchFhirCommandRequest requ
7476

7577
using (_distributedSemaphoreProvider.AcquireSemaphore(request.facilityId, request.queryConfig.MaxConcurrentRequests.GetValueOrDefault(1), _distributedLockSettings.Expiration, cancellationToken))
7678
{
79+
// Create a new handler chain using a DelegatingHandler around a base HttpClientHandler
80+
var innerHandler = new HttpClientHandler();
81+
var headerCapturingHandler = new HeaderCapturingHandler { InnerHandler = innerHandler };
82+
var httpClientWithHandler = new HttpClient(headerCapturingHandler);
7783

78-
var fhirClient = new FhirClient(request.queryConfig.FhirServerBaseUrl, _httpClient, new FhirClientSettings
84+
var fhirClient = new FhirClient(request.queryConfig.FhirServerBaseUrl, httpClientWithHandler, new FhirClientSettings
7985
{
8086
PreferredFormat = ResourceFormat.Json
8187
});
@@ -99,7 +105,12 @@ public async IAsyncEnumerable<Bundle> ExecuteAsync(SearchFhirCommandRequest requ
99105
resultBundle = await fhirClient.SearchAsync(request.searchParams, request.resourceType.ToString(), cancellationToken);
100106
}
101107
}
102-
catch(Exception ex)
108+
catch (FhirOperationException ex) when (ex.Status == HttpStatusCode.TooManyRequests)
109+
{
110+
var retryAfter = FhirCommandUtils.ParseRetryAfter(headerCapturingHandler.LastResponseHeaders);
111+
throw new TooManyRequestsException($"Too many requests for search on {request.resourceType}", retryAfter);
112+
}
113+
catch (Exception ex)
103114
{
104115
_logger.LogError(ex, "Error encountered while searching FHIR resources. ResourceType: {ResourceType}; FacilityId: {facilityId};", request.resourceType, request.facilityId.Sanitize());
105116
throw;
@@ -117,6 +128,11 @@ public async IAsyncEnumerable<Bundle> ExecuteAsync(SearchFhirCommandRequest requ
117128
{
118129
resultBundle = await fhirClient.ContinueAsync(resultBundle, ct: cancellationToken);
119130
}
131+
catch (FhirOperationException ex) when (ex.Status == HttpStatusCode.TooManyRequests)
132+
{
133+
var retryAfter = FhirCommandUtils.ParseRetryAfter(headerCapturingHandler.LastResponseHeaders);
134+
throw new TooManyRequestsException($"Too many requests during paging for {request.resourceType}", retryAfter);
135+
}
120136
catch (Exception ex)
121137
{
122138
_logger.LogError(ex, "Error encountered while searching FHIR resources. ResourceType: {ResourceType}; SearchParams: {SearchParams},\n\n\t{stack}\n\n\t{innerStack}", request.resourceType, request.searchParams, ex.StackTrace, ex.InnerException?.StackTrace);
@@ -146,7 +162,12 @@ public async Task<Bundle> ExecuteNonPagingAsync(SearchFhirCommandRequest request
146162

147163
using (_distributedSemaphoreProvider.AcquireSemaphore(request.facilityId, request.queryConfig.MaxConcurrentRequests.GetValueOrDefault(1), _distributedLockSettings.Expiration, cancellationToken))
148164
{
149-
var fhirClient = new FhirClient(request.queryConfig.FhirServerBaseUrl, _httpClient, new FhirClientSettings
165+
// Create a new handler chain using a DelegatingHandler around a base HttpClientHandler
166+
var innerHandler = new HttpClientHandler();
167+
var headerCapturingHandler = new HeaderCapturingHandler { InnerHandler = innerHandler };
168+
var httpClientWithHandler = new HttpClient(headerCapturingHandler);
169+
170+
var fhirClient = new FhirClient(request.queryConfig.FhirServerBaseUrl, httpClientWithHandler, new FhirClientSettings
150171
{
151172
PreferredFormat = ResourceFormat.Json
152173
});
@@ -157,7 +178,16 @@ public async Task<Bundle> ExecuteNonPagingAsync(SearchFhirCommandRequest request
157178
fhirClient.RequestHeaders.Authorization = (AuthenticationHeaderValue)authBuilderResults.authHeader;
158179
}
159180

160-
var resultBundle = await fhirClient.SearchAsync(request.searchParams, request.resourceType.ToString(), cancellationToken);
181+
Bundle resultBundle;
182+
try
183+
{
184+
resultBundle = await fhirClient.SearchAsync(request.searchParams, request.resourceType.ToString(), cancellationToken);
185+
}
186+
catch (FhirOperationException ex) when (ex.Status == HttpStatusCode.TooManyRequests)
187+
{
188+
var retryAfter = FhirCommandUtils.ParseRetryAfter(headerCapturingHandler.LastResponseHeaders);
189+
throw new TooManyRequestsException($"Too many requests for non-paging search on {request.resourceType}", retryAfter);
190+
}
161191
IncrementResourceAcquiredMetric(request.correlationId, request.patientId, request.facilityId, request.queryPhase.ToString(), request.resourceType.ToString(), resultBundle.Id);
162192
return resultBundle;
163193
}
@@ -174,4 +204,4 @@ private void IncrementResourceAcquiredMetric(string? correlationId, string? pati
174204
new KeyValuePair<string, object?>(DiagnosticNames.ResourceId, resourceId)
175205
]);
176206
}
177-
}
207+
}

0 commit comments

Comments
 (0)