shoc-backend/Api.SeaHavenIndustries/Infrastructure/ProcurementWorkOrderClient.cs
2026-07-24 22:13:25 -03:00

332 lines
14 KiB
C#

using System.Globalization;
using System.Net;
using System.Security.Cryptography;
using System.Text;
using System.Text.Json;
using System.Text.Json.Serialization;
using Amazon.Runtime;
using Microsoft.Extensions.Options;
using SeaHaven.Services.Configuration;
using SeaHaven.Services.Interfaces;
namespace Api.SeaHavenIndustries.Infrastructure
{
internal interface IProcurementAwsCredentialsProvider
{
Task<ImmutableCredentials> GetAsync(CancellationToken cancellationToken);
}
internal sealed class DefaultProcurementAwsCredentialsProvider
: IProcurementAwsCredentialsProvider
{
public async Task<ImmutableCredentials> GetAsync(CancellationToken cancellationToken)
{
var credentials = FallbackCredentialsFactory.GetCredentials();
return await credentials.GetCredentialsAsync().WaitAsync(cancellationToken);
}
}
internal sealed class ProcurementWorkOrderClient : IProcurementWorkOrderClient
{
private static readonly JsonSerializerOptions JsonOptions = new()
{
PropertyNameCaseInsensitive = false,
PropertyNamingPolicy = JsonNamingPolicy.SnakeCaseLower
};
private readonly HttpClient _httpClient;
private readonly IProcurementAwsCredentialsProvider _credentials;
private readonly IOptionsMonitor<WorkOrderReconciliationOptions> _options;
private readonly TimeProvider _timeProvider;
public ProcurementWorkOrderClient(
HttpClient httpClient,
IProcurementAwsCredentialsProvider credentials,
IOptionsMonitor<WorkOrderReconciliationOptions> options,
TimeProvider timeProvider)
{
_httpClient = httpClient;
_credentials = credentials;
_options = options;
_timeProvider = timeProvider;
}
public Task<ProcurementPage<ProcurementWorkOrder>> GetWorkOrdersAsync(
string? cursor,
int limit,
CancellationToken cancellationToken) =>
GetPageAsync<ProcurementWorkOrder>("/work-orders", cursor, limit, cancellationToken);
public Task<ProcurementPage<ProcurementWorkOrderComment>> GetCommentsAsync(
string workOrderId,
string? cursor,
int limit,
CancellationToken cancellationToken)
{
if (workOrderId.Length is 0 or > 64 || workOrderId.Any(c => c is < '0' or > '9'))
throw new ArgumentException("Work-order ID must be numeric.", nameof(workOrderId));
return GetPageAsync<ProcurementWorkOrderComment>(
$"/work-orders/{workOrderId}/comments",
cursor,
limit,
cancellationToken);
}
private async Task<ProcurementPage<T>> GetPageAsync<T>(
string path,
string? cursor,
int limit,
CancellationToken cancellationToken)
{
var options = _options.CurrentValue;
if (limit is < 1 or > 500)
throw new ArgumentOutOfRangeException(nameof(limit));
if (cursor is { Length: > 0 } && cursor.Length > options.MaxCursorLength)
throw new ArgumentException("Cursor is too long.", nameof(cursor));
var origin = ValidateOrigin(options.BaseUrl);
var query = new List<KeyValuePair<string, string>>
{
new("limit", limit.ToString(CultureInfo.InvariantCulture))
};
if (cursor != null)
query.Add(new("cursor", cursor));
var canonicalQuery = string.Join(
"&",
query.OrderBy(k => k.Key, StringComparer.Ordinal)
.Select(k => $"{Encode(k.Key)}={Encode(k.Value)}"));
var uri = new Uri(origin, $"{path}?{canonicalQuery}");
for (var attempt = 0; ; attempt++)
{
using var timeout = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
timeout.CancelAfter(TimeSpan.FromSeconds(options.RequestTimeoutSeconds));
HttpResponseMessage response;
try
{
using var request = new HttpRequestMessage(HttpMethod.Get, uri);
await SignAsync(request, path, canonicalQuery, options.Region, timeout.Token);
response = await _httpClient.SendAsync(
request,
HttpCompletionOption.ResponseHeadersRead,
timeout.Token);
}
catch (OperationCanceledException) when (
!cancellationToken.IsCancellationRequested
&& attempt < options.MaxRetries)
{
await DelayBeforeRetryAsync(attempt, options, cancellationToken);
continue;
}
catch (HttpRequestException) when (attempt < options.MaxRetries)
{
await DelayBeforeRetryAsync(attempt, options, cancellationToken);
continue;
}
using (response)
{
if ((response.StatusCode == HttpStatusCode.TooManyRequests
|| (int)response.StatusCode >= 500)
&& attempt < options.MaxRetries)
{
await DelayBeforeRetryAsync(attempt, options, cancellationToken);
continue;
}
if (!response.IsSuccessStatusCode)
throw new HttpRequestException(
$"Procurement API returned HTTP {(int)response.StatusCode}.",
null,
response.StatusCode);
var bytes = await ReadBoundedAsync(
response.Content,
options.MaxResponseBytes,
timeout.Token);
PageEnvelope<T>? envelope;
try
{
envelope = JsonSerializer.Deserialize<PageEnvelope<T>>(bytes, JsonOptions);
}
catch (JsonException ex)
{
throw new InvalidDataException("Procurement API returned invalid JSON.", ex);
}
if (envelope?.Items == null)
throw new InvalidDataException("Procurement API response is missing items.");
if (envelope.NextCursor is { Length: 0 }
|| envelope.NextCursor?.Length > options.MaxCursorLength)
{
throw new InvalidDataException("Procurement API response has an invalid cursor.");
}
ValidateItems(envelope.Items);
return new ProcurementPage<T>(envelope.Items, envelope.NextCursor);
}
}
}
private static void ValidateItems<T>(IEnumerable<T> items)
{
foreach (var item in items)
{
switch (item)
{
case ProcurementWorkOrder workOrder
when !IsNumericId(workOrder.WorkOrderId)
|| !ValidStatus(workOrder.WoStatus)
|| !ValidRecordType(workOrder.RecordType):
throw new InvalidDataException("Procurement API returned an invalid work order.");
case ProcurementWorkOrderComment comment
when !IsNumericId(comment.WorkOrderId)
|| string.IsNullOrWhiteSpace(comment.CommentId)
|| comment.CommentId.Length > 450:
throw new InvalidDataException("Procurement API returned an invalid comment.");
}
}
}
private static bool IsNumericId(string? value) =>
value is { Length: > 0 and <= 64 }
&& value.All(c => c is >= '0' and <= '9');
private static bool ValidStatus(string? value) =>
value == null
|| value is "new" or "assigned" or "in_progress" or "on_hold"
or "completed" or "cancelled" or "unknown";
private static bool ValidRecordType(string? value) =>
value == null
|| value is "new_work_order" or "update" or "comment" or "cancellation";
private async Task DelayBeforeRetryAsync(
int attempt,
WorkOrderReconciliationOptions options,
CancellationToken cancellationToken)
{
await Task.Delay(
TimeSpan.FromMilliseconds(
options.RetryBaseDelayMilliseconds * (1 << Math.Min(attempt, 8))),
_timeProvider,
cancellationToken);
}
private async Task SignAsync(
HttpRequestMessage request,
string path,
string canonicalQuery,
string region,
CancellationToken cancellationToken)
{
var credentials = await _credentials.GetAsync(cancellationToken);
var now = _timeProvider.GetUtcNow().UtcDateTime;
var amzDate = now.ToString("yyyyMMdd'T'HHmmss'Z'", CultureInfo.InvariantCulture);
var date = now.ToString("yyyyMMdd", CultureInfo.InvariantCulture);
const string payloadHash = "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855";
request.Headers.TryAddWithoutValidation("X-Amz-Date", amzDate);
request.Headers.TryAddWithoutValidation("X-Amz-Content-Sha256", payloadHash);
var signedHeaders = "host;x-amz-content-sha256;x-amz-date";
var canonicalHeaders = $"host:{request.RequestUri!.Host}\n"
+ $"x-amz-content-sha256:{payloadHash}\n"
+ $"x-amz-date:{amzDate}\n";
if (credentials.UseToken)
{
request.Headers.TryAddWithoutValidation("X-Amz-Security-Token", credentials.Token);
signedHeaders += ";x-amz-security-token";
canonicalHeaders += $"x-amz-security-token:{credentials.Token}\n";
}
var canonicalRequest = $"GET\n{path}\n{canonicalQuery}\n"
+ $"{canonicalHeaders}\n{signedHeaders}\n{payloadHash}";
var scope = $"{date}/{region}/execute-api/aws4_request";
var stringToSign = "AWS4-HMAC-SHA256\n"
+ $"{amzDate}\n{scope}\n{Sha256Hex(canonicalRequest)}";
var signingKey = DeriveSigningKey(credentials.SecretKey, date, region, "execute-api");
var signature = Convert.ToHexString(
HMACSHA256.HashData(signingKey, Encoding.UTF8.GetBytes(stringToSign)))
.ToLowerInvariant();
CryptographicOperations.ZeroMemory(signingKey);
request.Headers.TryAddWithoutValidation(
"Authorization",
$"AWS4-HMAC-SHA256 Credential={credentials.AccessKey}/{scope}, "
+ $"SignedHeaders={signedHeaders}, Signature={signature}");
}
private static Uri ValidateOrigin(string value)
{
if (!Uri.TryCreate(value, UriKind.Absolute, out var uri)
|| uri.Scheme != Uri.UriSchemeHttps
|| !string.Equals(
uri.Host,
"procurement-api.seahaven.com",
StringComparison.OrdinalIgnoreCase)
|| !uri.IsDefaultPort
|| uri.AbsolutePath != "/"
|| uri.Query.Length > 0
|| uri.UserInfo.Length > 0)
{
throw new InvalidOperationException("Procurement API base URL is not the allowed origin.");
}
return uri;
}
private static async Task<byte[]> ReadBoundedAsync(
HttpContent content,
int maximumBytes,
CancellationToken cancellationToken)
{
if (content.Headers.ContentLength > maximumBytes)
throw new InvalidDataException("Procurement API response is too large.");
await using var source = await content.ReadAsStreamAsync(cancellationToken);
using var destination = new MemoryStream();
var buffer = new byte[81920];
while (true)
{
var read = await source.ReadAsync(buffer, cancellationToken);
if (read == 0)
return destination.ToArray();
if (destination.Length + read > maximumBytes)
throw new InvalidDataException("Procurement API response is too large.");
destination.Write(buffer, 0, read);
}
}
private static byte[] DeriveSigningKey(
string secret,
string date,
string region,
string service)
{
var dateKey = HMACSHA256.HashData(
Encoding.UTF8.GetBytes($"AWS4{secret}"),
Encoding.UTF8.GetBytes(date));
var regionKey = HMACSHA256.HashData(dateKey, Encoding.UTF8.GetBytes(region));
CryptographicOperations.ZeroMemory(dateKey);
var serviceKey = HMACSHA256.HashData(regionKey, Encoding.UTF8.GetBytes(service));
CryptographicOperations.ZeroMemory(regionKey);
var signingKey = HMACSHA256.HashData(serviceKey, Encoding.UTF8.GetBytes("aws4_request"));
CryptographicOperations.ZeroMemory(serviceKey);
return signingKey;
}
private static string Sha256Hex(string value) =>
Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(value))).ToLowerInvariant();
private static string Encode(string value) =>
Uri.EscapeDataString(value).Replace("%7E", "~", StringComparison.Ordinal);
private sealed class PageEnvelope<T>
{
public List<T>? Items { get; set; }
[JsonPropertyName("next_cursor")]
public string? NextCursor { get; set; }
}
}
}