using System.Net; using System.Net.Http.Headers; using System.Security.Cryptography; using System.Text; using System.Text.Json; using FinalFactory.Rendezvous.Contracts; namespace FinalFactory.Rendezvous.Client; internal sealed class RendezvousHttpTransport { private readonly HttpClient _httpClient; private readonly RendezvousClientOptions _options; private readonly IRendezvousDelay _delay; internal RendezvousHttpTransport( HttpClient httpClient, RendezvousClientOptions? options, IRendezvousDelay? delay) { _httpClient = httpClient ?? throw new ArgumentNullException(nameof(httpClient)); RendezvousClientOptions suppliedOptions = options ?? new RendezvousClientOptions(); suppliedOptions.Validate(); _options = new RendezvousClientOptions { MaximumSafeRetries = suppliedOptions.MaximumSafeRetries, RequestTimeout = suppliedOptions.RequestTimeout, InitialRetryDelay = suppliedOptions.InitialRetryDelay, MaximumRetryDelay = suppliedOptions.MaximumRetryDelay, JitterRatio = suppliedOptions.JitterRatio, }; _delay = delay ?? new SystemRendezvousDelay(); } internal async Task> SendSafeAsync( Func requestFactory, CancellationToken cancellationToken) { for (int attempt = 0; ; attempt++) { cancellationToken.ThrowIfCancellationRequested(); using CancellationTokenSource requestTimeout = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); requestTimeout.CancelAfter(_options.RequestTimeout); CancellationToken requestCancellation = requestTimeout.Token; try { using HttpRequestMessage request = requestFactory(); using HttpResponseMessage response = await _httpClient .SendAsync(request, HttpCompletionOption.ResponseHeadersRead, requestCancellation) .ConfigureAwait(false); if (response.IsSuccessStatusCode) { if (typeof(T) == typeof(bool) && response.StatusCode == HttpStatusCode.NoContent) { return RendezvousClientResult.Success((T)(object)true); } byte[] payload; try { payload = await ReadBoundedAsync(response.Content, requestCancellation) .ConfigureAwait(false); } catch (InvalidDataException) { return RendezvousClientResult.Failure( RendezvousErrorCode.InternalError, "The service returned an oversized success response."); } T? value; try { value = JsonSerializer.Deserialize(payload, ContractJson.Options); } catch (JsonException) { value = default; } return value is null ? RendezvousClientResult.Failure( RendezvousErrorCode.InternalError, "The service returned an invalid success response.") : RendezvousClientResult.Success(value); } ApiError error = await ReadErrorAsync(response, requestCancellation).ConfigureAwait(false); int? retryAfter = error.RetryAfterSeconds ?? GetRetryAfterSeconds(response.Headers.RetryAfter); if (attempt < _options.MaximumSafeRetries && IsTransient(error.Code)) { await _delay.DelayAsync( GetRetryDelay(attempt, retryAfter), cancellationToken).ConfigureAwait(false); continue; } return RendezvousClientResult.Failure(error.Code, error.Message, retryAfter); } catch (Exception exception) when ( IsTransientTransportFailure(exception, cancellationToken) && attempt < _options.MaximumSafeRetries) { await _delay.DelayAsync(GetRetryDelay(attempt, null), cancellationToken) .ConfigureAwait(false); } catch (Exception exception) when (IsTransientTransportFailure(exception, cancellationToken)) { return RendezvousClientResult.Failure( RendezvousErrorCode.ServiceUnavailable, "The Rendezvous service did not return a valid response."); } } } internal static HttpRequestMessage JsonRequest( HttpMethod method, string uri, T body, string? publisherCredential = null) { HttpRequestMessage request = new(method, uri) { Content = new StringContent( JsonSerializer.Serialize(body, ContractJson.Options), Encoding.UTF8, "application/json"), }; if (publisherCredential is not null) { request.Headers.Authorization = new AuthenticationHeaderValue( "Bearer", RequireCredential(publisherCredential)); } return request; } internal static string RequireCredential(string credential) => !string.IsNullOrWhiteSpace(credential) ? credential : throw new ArgumentException("A publisher credential is required.", nameof(credential)); private static async Task ReadErrorAsync( HttpResponseMessage response, CancellationToken cancellationToken) { try { byte[] payload = await ReadBoundedAsync(response.Content, cancellationToken) .ConfigureAwait(false); ApiError? error = JsonSerializer.Deserialize(payload, ContractJson.Options); return error is not null && error.Code != RendezvousErrorCode.None ? error : FallbackError(response.StatusCode); } catch (Exception exception) when (exception is JsonException or InvalidDataException) { return FallbackError(response.StatusCode); } } private static async Task ReadBoundedAsync( HttpContent content, CancellationToken cancellationToken) { using Stream source = await content.ReadAsStreamAsync().ConfigureAwait(false); using MemoryStream destination = new(); byte[] buffer = new byte[8192]; while (true) { int read = await source.ReadAsync(buffer.AsMemory(), cancellationToken) .ConfigureAwait(false); if (read == 0) { return destination.ToArray(); } if (destination.Length + read > ContractLimits.BrowserResponseMaxBytes) { throw new InvalidDataException("The service response exceeded the SDK limit."); } await destination.WriteAsync(buffer.AsMemory(0, read), cancellationToken) .ConfigureAwait(false); } } private TimeSpan GetRetryDelay(int attempt, int? retryAfterSeconds) { TimeSpan basis = retryAfterSeconds.HasValue ? TimeSpan.FromSeconds(Math.Max(0, retryAfterSeconds.Value)) : TimeSpan.FromMilliseconds( _options.InitialRetryDelay.TotalMilliseconds * Math.Pow(2, attempt)); double bounded = Math.Min(basis.TotalMilliseconds, _options.MaximumRetryDelay.TotalMilliseconds); if (_options.JitterRatio == 0 || bounded == 0) { return TimeSpan.FromMilliseconds(bounded); } byte[] random = new byte[1]; RandomNumberGenerator.Fill(random); double unit = random[0] / 255d; double multiplier = 1 - _options.JitterRatio + (2 * _options.JitterRatio * unit); return TimeSpan.FromMilliseconds(Math.Min( bounded * multiplier, _options.MaximumRetryDelay.TotalMilliseconds)); } private static bool IsTransient(RendezvousErrorCode code) => code is RendezvousErrorCode.RateLimited or RendezvousErrorCode.CapacityExceeded or RendezvousErrorCode.ServiceUnavailable; private static bool IsTransientTransportFailure( Exception exception, CancellationToken callerCancellation) => exception is HttpRequestException or IOException || exception is OperationCanceledException && !callerCancellation.IsCancellationRequested; private static int? GetRetryAfterSeconds(RetryConditionHeaderValue? retryAfter) => retryAfter?.Delta is TimeSpan delta ? Math.Max(0, (int)Math.Ceiling(delta.TotalSeconds)) : null; private static ApiError FallbackError(HttpStatusCode statusCode) => new() { Code = statusCode switch { HttpStatusCode.BadRequest => RendezvousErrorCode.InvalidRequest, HttpStatusCode.Unauthorized => RendezvousErrorCode.AuthenticationRequired, HttpStatusCode.Forbidden => RendezvousErrorCode.Forbidden, HttpStatusCode.NotFound => RendezvousErrorCode.NotFound, HttpStatusCode.Conflict => RendezvousErrorCode.Conflict, HttpStatusCode.Gone => RendezvousErrorCode.Expired, HttpStatusCode.TooManyRequests => RendezvousErrorCode.RateLimited, HttpStatusCode.RequestTimeout => RendezvousErrorCode.ServiceUnavailable, HttpStatusCode.BadGateway => RendezvousErrorCode.ServiceUnavailable, HttpStatusCode.ServiceUnavailable => RendezvousErrorCode.ServiceUnavailable, HttpStatusCode.GatewayTimeout => RendezvousErrorCode.ServiceUnavailable, _ => RendezvousErrorCode.InternalError, }, Message = "The service returned an error without a valid Rendezvous envelope.", }; }