Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -1,13 +1,13 @@
using System.Diagnostics;
using DataAcquisition.Domain.Application.Models;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Models.Api.QueryLog;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Models.Exceptions;
using LantanaGroup.Link.DataAcquisition.Domain.Application.Queries;
using LantanaGroup.Link.DataAcquisition.Domain.Infrastructure;
using LantanaGroup.Link.DataAcquisition.Domain.Infrastructure.Entities;
using LantanaGroup.Link.DataAcquisition.Domain.Infrastructure.Models.Enums;
using LantanaGroup.Link.DataAcquisition.Domain.Models;
using LantanaGroup.Link.Shared.Application.Models.Telemetry;
using Microsoft.Extensions.Logging;
using System.Diagnostics;

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

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

public class DataAcquisitionLogManager : IDataAcquisitionLogManager
{
public readonly ILogger<DataAcquisitionLogManager> _logger;
public readonly IDatabase _database;
private readonly IDataAcquisitionLogQueries _logQueries;

public DataAcquisitionLogManager(ILogger<DataAcquisitionLogManager> logger, IDatabase database)
public DataAcquisitionLogManager(ILogger<DataAcquisitionLogManager> logger, IDatabase database, IDataAcquisitionLogQueries logQueries)
{
_logger = logger ?? throw new ArgumentNullException(nameof(logger));
_database = database ?? throw new ArgumentNullException(nameof(database));
_logQueries = logQueries;
}

public async Task<DataAcquisitionLogModel> CreateAsync(CreateDataAcquisitionLogModel model, CancellationToken cancellationToken = default)
Expand Down Expand Up @@ -202,4 +205,30 @@ public async Task UpdateTailFlagForFacilityCorrelationIdReportTrackingId(List<lo
await _database.DataAcquisitionLogRepository.SaveChangesAsync();
}
}

public async Task ThrottleFacilityAcquisitions(string facilityId, DateTime executionDate, CancellationToken cancellationToken = default)
{
long? lastId = null;
var batchSize = 1000;
while (true)
{
// Get All Active Logs For Batch
var toThrottle = await _logQueries.GetNextEligibleBatchForFacility(facilityId, lastId, batchSize, [RequestStatus.Failed, RequestStatus.Ready, RequestStatus.Pending], executionDate, cancellationToken);

if(toThrottle.Count == 0)
{
break;
}

//Update their next processing time
foreach (var log in toThrottle)
{
log.ExecutionDate = executionDate;
}

await _database.DataAcquisitionLogRepository.SaveChangesAsync();

lastId = toThrottle.Max(l => l.Id);
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
namespace LantanaGroup.Link.DataAcquisition.Domain.Application.Models.Exceptions
{
public class ProcessingDelayException : Exception
{
public ProcessingDelayException(string message) : base(message) { }
}
}
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) : base(message, System.Net.HttpStatusCode.TooManyRequests)
{
RetryAfter = retryAfter;
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ public interface IDataAcquisitionLogQueries

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

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

Task<List<string>> GetResourceIdsForReportPatient(string correlationId, string facilityId, string resourceType, CancellationToken cancellationToken = default);
}
Expand Down Expand Up @@ -566,13 +566,16 @@ public async Task<List<string>> GetFacilitiesWithPendingAndRetryableFailedReques
.ToListAsync(cancellationToken);
}

public async Task<List<DataAcquisitionLogModel>> GetNextEligibleBatchForFacility(string facilityId, long? lastId, int batchSize, CancellationToken cancellationToken = default)
public async Task<List<DataAcquisitionLogModel>> GetNextEligibleBatchForFacility(string facilityId, long? lastId, int batchSize, List<RequestStatus> statuses, DateTime? designagtedExecutionTime = null, CancellationToken cancellationToken = default)
{
designagtedExecutionTime ??= DateTime.UtcNow;

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.ExecutionDate == null || log.ExecutionDate <= designagtedExecutionTime)
&& (log.Status == null || statuses.Contains(log.Status.Value))
orderby log.Id
select new DataAcquisitionLogModel
{
Id = log.Id,
Expand Down Expand Up @@ -622,4 +625,9 @@ private Expression<Func<T, object>> SetSortBy<T>(string? sortBy)
var converted = Expression.Convert(property, typeof(object));
return Expression.Lambda<Func<T, object>>(converted, parameter);
}

public Task<List<DataAcquisitionLogModel>> GetNextBatchForFacility(string facilityId, long? lastId, int batchSize, CancellationToken cancellationToken = default)
{
throw new NotImplementedException();
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
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;
}
}

public static class FhirCommandUtils
{
public static TimeSpan ParseRetryAfter(HttpResponseHeaders? headers, TimeSpan defaultDelay = default)
{
defaultDelay = defaultDelay == default ? TimeSpan.FromSeconds(60) : defaultDelay;

if (headers == null || headers.RetryAfter == null)
{
return defaultDelay;
}

var retryValue = headers.RetryAfter;
TimeSpan delay;

if (retryValue.Delta.HasValue)
{
delay = retryValue.Delta.Value;
}
else if (retryValue.Date.HasValue)
{
delay = retryValue.Date.Value - DateTimeOffset.UtcNow;
}
else
{
return defaultDelay;
}

if (delay <= TimeSpan.Zero)
{
return defaultDelay;
}

return delay;
}
}
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)
]);
}
}
}
Loading
Loading