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 GetAsync(CancellationToken cancellationToken); } internal sealed class DefaultProcurementAwsCredentialsProvider : IProcurementAwsCredentialsProvider { public async Task 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 _options; private readonly TimeProvider _timeProvider; public ProcurementWorkOrderClient( HttpClient httpClient, IProcurementAwsCredentialsProvider credentials, IOptionsMonitor options, TimeProvider timeProvider) { _httpClient = httpClient; _credentials = credentials; _options = options; _timeProvider = timeProvider; } public Task> GetWorkOrdersAsync( string? cursor, int limit, CancellationToken cancellationToken) => GetPageAsync("/work-orders", cursor, limit, cancellationToken); public Task> 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( $"/work-orders/{workOrderId}/comments", cursor, limit, cancellationToken); } private async Task> GetPageAsync( 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> { 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? envelope; try { envelope = JsonSerializer.Deserialize>(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(envelope.Items, envelope.NextCursor); } } } private static void ValidateItems(IEnumerable 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 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 { public List? Items { get; set; } [JsonPropertyName("next_cursor")] public string? NextCursor { get; set; } } } }