Skip to content
Merged
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
using System;
using Hl7.Fhir.Rest;

namespace LantanaGroup.Link.DataAcquisition.Domain.Application.Models.Exceptions
{
public class TooManyRequestsException : FhirOperationException
{
public TimeSpan? RetryAfter { get; }
public TooManyRequestsException(string message, TimeSpan? retryAfter = null) : base(message, System.Net.HttpStatusCode.TooManyRequests)
{
RetryAfter = retryAfter;
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -569,10 +569,11 @@ public async Task<List<string>> GetFacilitiesWithPendingAndRetryableFailedReques
public async Task<List<DataAcquisitionLogModel>> GetNextEligibleBatchForFacility(string facilityId, long? lastId, int batchSize, CancellationToken cancellationToken = default)
{
var query = from log in _dbContext.DataAcquisitionLogs
orderby log.Id
where log.FacilityId == facilityId
&& (lastId == null || log.Id > lastId)
&& (log.Status == RequestStatus.Pending || log.Status == RequestStatus.Failed)
&& (log.Status == RequestStatus.Pending ||
(log.Status == RequestStatus.Failed && log.RetryAttempts < DataAcquisitionLog.MaxRetryAttempts))
orderby log.Priority descending, log.ExecutionDate ascending, log.Id ascending
Comment thread
nvmLantana marked this conversation as resolved.
Outdated
select new DataAcquisitionLogModel
{
Id = log.Id,
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
using System.Net.Http.Headers;

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

internal class HeaderCapturingHandler : DelegatingHandler
{
public HttpResponseHeaders? LastResponseHeaders { get; private set; }

protected override async Task<HttpResponseMessage> SendAsync(HttpRequestMessage request, CancellationToken cancellationToken)
{
var response = await base.SendAsync(request, cancellationToken);
LastResponseHeaders = response.Headers;
return response;
}
}

internal static class FhirCommandUtils
{
public static TimeSpan? ParseRetryAfter(HttpResponseHeaders? headers)
{
if (headers == null) return null;

if (headers.TryGetValues("Retry-After", out var values))
{
var value = values.FirstOrDefault();
if (int.TryParse(value, out int seconds))
return TimeSpan.FromSeconds(seconds);
if (DateTimeOffset.TryParse(value, out var date))
return date - DateTimeOffset.UtcNow;
}
Comment thread
nvmLantana marked this conversation as resolved.
Outdated
return null;
}
}
Original file line number Diff line number Diff line change
@@ -1,13 +1,15 @@
using Hl7.Fhir.Model;
using Hl7.Fhir.Rest;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Factories.Auth;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Interfaces;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Models;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Models.Exceptions;
using LantanaGroup.Link.Shared.Application.Models.Configs;
using Medallion.Threading;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using System.Net;
using System.Net.Http.Headers;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Factories.Auth;

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

Expand All @@ -18,12 +20,13 @@ public record ReadFhirCommandRequest(
string baseUrl,
FhirQueryConfigurationModel fhirQueryConfiguration);

public interface IReadFhirCommand
{
public interface IReadFhirCommand
{
Task<DomainResource> ExecuteAsync(
ReadFhirCommandRequest request,
CancellationToken cancellationToken = default);
}

public class ReadFhirCommand : IReadFhirCommand
{
private readonly ILogger<ReadFhirCommand> _logger;
Expand All @@ -48,21 +51,23 @@ public ReadFhirCommand(

public async Task<DomainResource> ExecuteAsync(ReadFhirCommandRequest request, CancellationToken cancellationToken = default)
{

if (string.IsNullOrWhiteSpace(request.resourceId))
throw new ArgumentNullException(nameof(request.resourceId), "Resource ID cannot be null or empty.");

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


if (request.fhirQueryConfiguration == null)
throw new ArgumentNullException(nameof(request.fhirQueryConfiguration), "FhirQueryConfiguration cannot be null.");


using (_distributedSemaphoreProvider.AcquireSemaphore(request.facilityId, request.fhirQueryConfiguration.MaxConcurrentRequests.GetValueOrDefault(1), _distributedLockSettings.Expiration, cancellationToken))
{
var fhirClient = new FhirClient(request.baseUrl.Trim('/'), _httpClient, new FhirClientSettings
// Create a new handler chain using a DelegatingHandler around a base HttpClientHandler
var innerHandler = new HttpClientHandler();
var headerCapturingHandler = new HeaderCapturingHandler { InnerHandler = innerHandler };
var httpClientWithHandler = new HttpClient(headerCapturingHandler);

var fhirClient = new FhirClient(request.baseUrl.Trim('/'), httpClientWithHandler, new FhirClientSettings
Comment thread
nvmLantana marked this conversation as resolved.
{
PreferredFormat = ResourceFormat.Json
});
Expand All @@ -73,15 +78,23 @@ public async Task<DomainResource> ExecuteAsync(ReadFhirCommandRequest request, C
fhirClient.RequestHeaders.Authorization = (AuthenticationHeaderValue)authBuilderResults.authHeader;
}


string location = request.resourceType switch
{
ResourceType.List => $"List/{request.resourceId}",
//ResourceType.Patient => TEMPORARYPatientIdPart(id),
_ => $"{request.resourceType}/{request.resourceId}"
};

var readResource = await fhirClient.ReadAsync<DomainResource>(location);
DomainResource readResource;
try
{
readResource = await fhirClient.ReadAsync<DomainResource>(location);
}
catch (FhirOperationException ex) when (ex.Status == HttpStatusCode.TooManyRequests)
{
var retryAfter = FhirCommandUtils.ParseRetryAfter(headerCapturingHandler.LastResponseHeaders);
throw new TooManyRequestsException($"Too many requests for {location}", retryAfter);
}

if (readResource == null)
{
Expand All @@ -91,4 +104,4 @@ public async Task<DomainResource> ExecuteAsync(ReadFhirCommandRequest request, C
return readResource;
}
}
}
}
Original file line number Diff line number Diff line change
@@ -1,16 +1,18 @@
using Hl7.Fhir.Model;
using Hl7.Fhir.Rest;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Factories.Auth;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Interfaces;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Models;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Models.Exceptions;
using LantanaGroup.Link.DataAcquisition.Domain.Infrastructure.Models.Enums;
using LantanaGroup.Link.Shared.Application.Models.Configs;
using LantanaGroup.Link.Shared.Application.Models.Telemetry;
using LantanaGroup.Link.Shared.Application.Services.Security;
using Medallion.Threading;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using System.Net;
using System.Net.Http.Headers;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Factories.Auth;
using ResourceType = Hl7.Fhir.Model.ResourceType;

namespace LantanaGroup.Link.DataAcquisition.Domain.Application.Services.FhirApi.Commands;
Expand All @@ -22,7 +24,7 @@ public record SearchFhirCommandRequest(
string? facilityId,
string? patientId,
string? correlationId,
QueryPhase? queryPhase,
QueryPhase? queryPhase,
FhirQueryType queryType
);

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

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

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

var fhirClient = new FhirClient(request.queryConfig.FhirServerBaseUrl, _httpClient, new FhirClientSettings
var fhirClient = new FhirClient(request.queryConfig.FhirServerBaseUrl, httpClientWithHandler, new FhirClientSettings
Comment thread
nvmLantana marked this conversation as resolved.
{
PreferredFormat = ResourceFormat.Json
});
Expand All @@ -99,7 +105,12 @@ public async IAsyncEnumerable<Bundle> ExecuteAsync(SearchFhirCommandRequest requ
resultBundle = await fhirClient.SearchAsync(request.searchParams, request.resourceType.ToString(), cancellationToken);
}
}
catch(Exception ex)
catch (FhirOperationException ex) when (ex.Status == HttpStatusCode.TooManyRequests)
{
var retryAfter = FhirCommandUtils.ParseRetryAfter(headerCapturingHandler.LastResponseHeaders);
throw new TooManyRequestsException($"Too many requests for search on {request.resourceType}", retryAfter);
}
catch (Exception ex)
{
_logger.LogError(ex, "Error encountered while searching FHIR resources. ResourceType: {ResourceType}; FacilityId: {facilityId};", request.resourceType, request.facilityId.Sanitize());
throw;
Expand All @@ -117,6 +128,11 @@ public async IAsyncEnumerable<Bundle> ExecuteAsync(SearchFhirCommandRequest requ
{
resultBundle = await fhirClient.ContinueAsync(resultBundle, ct: cancellationToken);
}
catch (FhirOperationException ex) when (ex.Status == HttpStatusCode.TooManyRequests)
{
var retryAfter = FhirCommandUtils.ParseRetryAfter(headerCapturingHandler.LastResponseHeaders);
throw new TooManyRequestsException($"Too many requests during paging for {request.resourceType}", retryAfter);
}
catch (Exception ex)
{
_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);
Expand Down Expand Up @@ -146,7 +162,12 @@ public async Task<Bundle> ExecuteNonPagingAsync(SearchFhirCommandRequest request

using (_distributedSemaphoreProvider.AcquireSemaphore(request.facilityId, request.queryConfig.MaxConcurrentRequests.GetValueOrDefault(1), _distributedLockSettings.Expiration, cancellationToken))
{
var fhirClient = new FhirClient(request.queryConfig.FhirServerBaseUrl, _httpClient, new FhirClientSettings
// Create a new handler chain using a DelegatingHandler around a base HttpClientHandler
var innerHandler = new HttpClientHandler();
var headerCapturingHandler = new HeaderCapturingHandler { InnerHandler = innerHandler };
var httpClientWithHandler = new HttpClient(headerCapturingHandler);

var fhirClient = new FhirClient(request.queryConfig.FhirServerBaseUrl, httpClientWithHandler, new FhirClientSettings
{
PreferredFormat = ResourceFormat.Json
});
Expand All @@ -157,7 +178,16 @@ public async Task<Bundle> ExecuteNonPagingAsync(SearchFhirCommandRequest request
fhirClient.RequestHeaders.Authorization = (AuthenticationHeaderValue)authBuilderResults.authHeader;
}

var resultBundle = await fhirClient.SearchAsync(request.searchParams, request.resourceType.ToString(), cancellationToken);
Bundle resultBundle;
try
{
resultBundle = await fhirClient.SearchAsync(request.searchParams, request.resourceType.ToString(), cancellationToken);
}
catch (FhirOperationException ex) when (ex.Status == HttpStatusCode.TooManyRequests)
{
var retryAfter = FhirCommandUtils.ParseRetryAfter(headerCapturingHandler.LastResponseHeaders);
throw new TooManyRequestsException($"Too many requests for non-paging search on {request.resourceType}", retryAfter);
}
IncrementResourceAcquiredMetric(request.correlationId, request.patientId, request.facilityId, request.queryPhase.ToString(), request.resourceType.ToString(), resultBundle.Id);
return resultBundle;
}
Expand All @@ -174,4 +204,4 @@ private void IncrementResourceAcquiredMetric(string? correlationId, string? pati
new KeyValuePair<string, object?>(DiagnosticNames.ResourceId, resourceId)
]);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -2,14 +2,13 @@
using DataAcquisition.Domain.Application.Models;
using Hl7.Fhir.Model;
using Hl7.Fhir.Rest;
using Hl7.Fhir.Serialization;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Managers;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Models;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Models.Exceptions;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Models.Factory;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Models.Kafka;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Queries;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Services.FhirApi.Commands;
using LantanaGroup.Link.DataAcquisition.Domain.Infrastructure.Entities;
using LantanaGroup.Link.DataAcquisition.Domain.Infrastructure.Models.Enums;
using LantanaGroup.Link.DataAcquisition.Domain.Models;
using LantanaGroup.Link.DataAcquisition.Domain.Settings;
Expand Down Expand Up @@ -115,6 +114,10 @@ await GenerateResourceAcquiredMessage(new ResourceAcquired

return resourceIds;
}
catch (TooManyRequestsException ex)
{
throw; // Propagate to higher level
}
catch (FhirOperationException ex)
{
if (fhirQuery.IsReference.GetValueOrDefault() && (ex.Status == HttpStatusCode.NotFound || ex.Status == HttpStatusCode.Gone))
Expand Down Expand Up @@ -190,7 +193,7 @@ private async Task<List<string>> ExecutePagingSearch(DataAcquisitionLogModel log

foreach (var resource in resources)
{
if(fhirQuery.IsReference.HasValue && fhirQuery.IsReference.Value)
if (fhirQuery.IsReference.HasValue && fhirQuery.IsReference.Value)
{
//if this is a reference resource, we need to handle it differently
await HandleReferenceResource(log, resource, cancellationToken);
Expand All @@ -211,6 +214,10 @@ await GenerateResourceAcquiredMessage(new ResourceAcquired

return resourceIds;
}
catch (TooManyRequestsException ex)
{
throw; // Propagate to higher level
}
catch (FhirOperationException ex)
{
if (ex.Outcome != null)
Expand Down Expand Up @@ -293,22 +300,22 @@ await _kafkaProducer.ProduceAsync(
_kafkaProducer.Flush(cancellationToken);
}

private void InsertDateExtension(DomainResource resource)
private void InsertDateExtension(DomainResource resource)
{
if(resource == null)
if (resource == null)
throw new ArgumentNullException(nameof(resource));

if(resource.Meta == null)
if (resource.Meta == null)
{
resource.Meta = new Meta();
resource.Meta.Extension = new List<Extension> { };
}

if(resource.Meta.Extension == null)
if (resource.Meta.Extension == null)
resource.Meta.Extension = new List<Extension> { };

if (!resource.Extension.Any(e => e.Url == DataAcquisitionConstants.Extension.DateReceivedExtensionUri))
resource.Meta.Extension.Add(new Extension { Url = DataAcquisitionConstants.Extension.DateReceivedExtensionUri, Value = new FhirDateTime(DateTime.UtcNow.ToString("yyyy-MM-ddTHH:mm:ssZ"))});
resource.Meta.Extension.Add(new Extension { Url = DataAcquisitionConstants.Extension.DateReceivedExtensionUri, Value = new FhirDateTime(DateTime.UtcNow.ToString("yyyy-MM-ddTHH:mm:ssZ")) });
}
#endregion
}
}
Loading
Loading