From 222705acdfc24e2bd7d570bcdcbe6485bef5a472 Mon Sep 17 00:00:00 2001 From: Michael Staib Date: Fri, 31 Jul 2026 15:09:26 +0200 Subject: [PATCH 1/2] Add Nitro seed update monitoring to Fusion Aspire --- .../GraphQLResourceBuilderExtensions.cs | 2 + .../HotChocolate.Fusion.Aspire.csproj | 9 + .../Nitro/INitroStageUpdateClient.cs | 24 + .../Nitro/NitroCompositionOptions.cs | 5 + .../Fusion.Aspire/Nitro/NitroGatewaySeed.cs | 6 +- .../Nitro/NitroOperationDocuments.cs | 22 + .../Nitro/NitroSeedCoordinator.cs | 337 ++++++++++- .../Nitro/NitroSeedUpdateMonitor.cs | 481 +++++++++++++++ .../Nitro/NitroSeedUpdateNotifier.cs | 54 ++ .../Nitro/NitroSeedUpdateOptions.cs | 21 + .../Nitro/NitroSeedUpdateService.cs | 201 +++++++ .../Nitro/NitroSeedUpdateState.cs | 34 ++ .../Fusion.Aspire/Nitro/NitroStageSnapshot.cs | 116 ++++ .../Nitro/NitroStageUpdateClient.cs | 272 +++++++++ .../Operations/GetNitroStageVersion.graphql | 19 + .../GetNitroStageVersion.graphql.sha256 | 1 + .../Nitro/Operations/WatchNitroStage.graphql | 34 ++ .../Operations/WatchNitroStage.graphql.sha256 | 1 + .../src/Fusion.Aspire/NitroExtensions.cs | 130 +++- .../src/Fusion.Aspire/SchemaComposition.cs | 199 ++++++- .../SchemaCompositionRegistration.cs | 3 + .../Fusion.Aspire.Tests/CompositionHarness.cs | 22 + .../Nitro/NitroOperationDocumentsTests.cs | 36 ++ .../Nitro/NitroSchemaCompositionTests.cs | 146 ++++- .../Nitro/NitroSeedUpdateMonitorTests.cs | 556 ++++++++++++++++++ .../Nitro/NitroStageUpdateClientTests.cs | 144 +++++ .../NitroExtensionsTests.cs | 162 +++++ 27 files changed, 3016 insertions(+), 21 deletions(-) create mode 100644 src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/INitroStageUpdateClient.cs create mode 100644 src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateMonitor.cs create mode 100644 src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateNotifier.cs create mode 100644 src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateOptions.cs create mode 100644 src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateService.cs create mode 100644 src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateState.cs create mode 100644 src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroStageSnapshot.cs create mode 100644 src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroStageUpdateClient.cs create mode 100644 src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/Operations/GetNitroStageVersion.graphql create mode 100644 src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/Operations/GetNitroStageVersion.graphql.sha256 create mode 100644 src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/Operations/WatchNitroStage.graphql create mode 100644 src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/Operations/WatchNitroStage.graphql.sha256 create mode 100644 src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroSeedUpdateMonitorTests.cs create mode 100644 src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroStageUpdateClientTests.cs diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/GraphQLResourceBuilderExtensions.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/GraphQLResourceBuilderExtensions.cs index 12856bd9890..faa06156568 100644 --- a/src/HotChocolate/Fusion/src/Fusion.Aspire/GraphQLResourceBuilderExtensions.cs +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/GraphQLResourceBuilderExtensions.cs @@ -94,6 +94,8 @@ public static IResourceBuilder WithGraphQLSchemaComposition( Settings = settings }); + NitroExtensions.TryAddAutoUpdateCommands(builder); + if (!builder.Resource.Annotations .OfType() .Any(command => command.Name == "recompose")) diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/HotChocolate.Fusion.Aspire.csproj b/src/HotChocolate/Fusion/src/Fusion.Aspire/HotChocolate.Fusion.Aspire.csproj index 1a886f4bab5..bee2f052a09 100644 --- a/src/HotChocolate/Fusion/src/Fusion.Aspire/HotChocolate.Fusion.Aspire.csproj +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/HotChocolate.Fusion.Aspire.csproj @@ -21,6 +21,7 @@ + @@ -34,18 +35,26 @@ + + + + + + + + diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/INitroStageUpdateClient.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/INitroStageUpdateClient.cs new file mode 100644 index 00000000000..d59f726e1c5 --- /dev/null +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/INitroStageUpdateClient.cs @@ -0,0 +1,24 @@ +namespace HotChocolate.Fusion.Aspire.Nitro; + +internal interface INitroStageUpdateClient +{ + Task SubscribeAsync( + NitroConnection connection, + string apiId, + string stage, + CancellationToken cancellationToken); + + Task GetLatestSnapshotAsync( + NitroConnection connection, + string apiId, + string stage, + CancellationToken cancellationToken); +} + +internal abstract class NitroStageSubscription : IAsyncDisposable +{ + public abstract IAsyncEnumerable ReadChangesAsync( + CancellationToken cancellationToken); + + public abstract ValueTask DisposeAsync(); +} diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroCompositionOptions.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroCompositionOptions.cs index b3317749124..2756a40a412 100644 --- a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroCompositionOptions.cs +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroCompositionOptions.cs @@ -16,4 +16,9 @@ internal sealed class NitroCompositionOptions /// Gets or sets the caller-supplied Nitro portal URL. /// public Uri? PortalUrl { get; set; } + + /// + /// Gets the options that control stage update detection and automatic adoption. + /// + public NitroSeedUpdateOptions SeedUpdates { get; } = new(); } diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroGatewaySeed.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroGatewaySeed.cs index 4994f6ec339..c30f9a4b26e 100644 --- a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroGatewaySeed.cs +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroGatewaySeed.cs @@ -22,12 +22,16 @@ namespace HotChocolate.Fusion.Aspire.Nitro; /// /// Whether the fusion configuration was downloaded for this run instead of taken from the cache. /// +/// +/// The SHA-256 hash of the gateway schema carried by the configuration. +/// internal sealed record NitroGatewaySeed( string ApiId, string Stage, string FilePath, DateTimeOffset DownloadedAt, - bool IsFresh); + bool IsFresh, + string SchemaHash); /// /// The outcome of acquiring the fusion configuration of a gateway. diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroOperationDocuments.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroOperationDocuments.cs index c94089384bb..e1d47586579 100644 --- a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroOperationDocuments.cs +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroOperationDocuments.cs @@ -14,18 +14,26 @@ internal static class NitroOperationDocuments private const string ResolveApiNameHashFile = "ResolveNitroApiName.graphql.sha256"; private const string ValidateSchemaHashFile = "ValidateNitroSchema.graphql.sha256"; private const string PollSchemaValidationHashFile = "PollNitroSchemaValidation.graphql.sha256"; + private const string GetStageVersionHashFile = "GetNitroStageVersion.graphql.sha256"; + private const string WatchStageHashFile = "WatchNitroStage.graphql.sha256"; private static string? s_resolveApiNameHash; private static string? s_validateSchemaHash; private static string? s_pollSchemaValidationHash; + private static string? s_getStageVersionHash; + private static string? s_watchStageHash; #else private const string ResolveApiNameFile = "ResolveNitroApiName.graphql"; private const string ValidateSchemaFile = "ValidateNitroSchema.graphql"; private const string PollSchemaValidationFile = "PollNitroSchemaValidation.graphql"; + private const string GetStageVersionFile = "GetNitroStageVersion.graphql"; + private const string WatchStageFile = "WatchNitroStage.graphql"; private static string? s_resolveApiName; private static string? s_validateSchema; private static string? s_pollSchemaValidation; + private static string? s_getStageVersion; + private static string? s_watchStage; #endif /// @@ -34,6 +42,8 @@ internal static class NitroOperationDocuments public const string ResolveApiNameOperationName = "ResolveNitroApiName"; public const string ValidateSchemaOperationName = "ValidateNitroSchema"; public const string PollSchemaValidationOperationName = "PollNitroSchemaValidation"; + public const string GetStageVersionOperationName = "GetNitroStageVersion"; + public const string WatchStageOperationName = "WatchNitroStage"; #if NITRO_PERSISTED_OPERATIONS /// @@ -47,6 +57,12 @@ public static string GetValidateSchemaOperationId() public static string GetPollSchemaValidationOperationId() => s_pollSchemaValidationHash ??= ReadDocument(PollSchemaValidationHashFile).Trim(); + + public static string GetStageVersionOperationId() + => s_getStageVersionHash ??= ReadDocument(GetStageVersionHashFile).Trim(); + + public static string GetWatchStageOperationId() + => s_watchStageHash ??= ReadDocument(WatchStageHashFile).Trim(); #else /// /// Gets the document that resolves the name of an api by its id. @@ -59,6 +75,12 @@ public static string GetValidateSchemaDocument() public static string GetPollSchemaValidationDocument() => s_pollSchemaValidation ??= ReadDocument(PollSchemaValidationFile); + + public static string GetStageVersionDocument() + => s_getStageVersion ??= ReadDocument(GetStageVersionFile); + + public static string GetWatchStageDocument() + => s_watchStage ??= ReadDocument(WatchStageFile); #endif private static string ReadDocument(string fileName) diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedCoordinator.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedCoordinator.cs index 2a74be20b49..6716b2af06f 100644 --- a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedCoordinator.cs +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedCoordinator.cs @@ -1,3 +1,5 @@ +using System.Security.Cryptography; +using HotChocolate.Fusion.Packaging; using HotChocolate.Transport.Http; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Logging.Abstractions; @@ -12,12 +14,15 @@ namespace HotChocolate.Fusion.Aspire.Nitro; /// internal sealed class NitroSeedCoordinator { - private readonly Dictionary _seedsByGateway = [with(StringComparer.Ordinal)]; + private readonly Dictionary _statesByGateway = [with(StringComparer.Ordinal)]; private readonly Lock _sync = new(); private readonly NitroConnectionResolver _connectionResolver; private readonly NitroSeedProvider _seedProvider; private readonly INitroSchemaValidator _schemaValidator; + private readonly INitroStageUpdateClient _stageUpdateClient; private readonly string _runSeedDirectory; + private readonly bool _initialAutoUpdate; + private long _nextRunSeedId; /// /// Initializes a new instance of . @@ -34,27 +39,38 @@ internal sealed class NitroSeedCoordinator /// /// The client that validates composed gateway schemas against Nitro. /// + /// + /// The client that observes the current version of the Nitro stage. + /// /// /// The directory that holds the private copies of this run. /// + /// + /// Whether newly observed configurations are applied automatically by default. + /// public NitroSeedCoordinator( string stage, NitroConnectionResolver connectionResolver, NitroSeedProvider seedProvider, INitroSchemaValidator schemaValidator, - string runSeedDirectory) + INitroStageUpdateClient stageUpdateClient, + string runSeedDirectory, + bool initialAutoUpdate) { ArgumentException.ThrowIfNullOrWhiteSpace(stage); ArgumentNullException.ThrowIfNull(connectionResolver); ArgumentNullException.ThrowIfNull(seedProvider); ArgumentNullException.ThrowIfNull(schemaValidator); + ArgumentNullException.ThrowIfNull(stageUpdateClient); ArgumentException.ThrowIfNullOrWhiteSpace(runSeedDirectory); Stage = stage; _connectionResolver = connectionResolver; _seedProvider = seedProvider; _schemaValidator = schemaValidator; + _stageUpdateClient = stageUpdateClient; _runSeedDirectory = runSeedDirectory; + _initialAutoUpdate = initialAutoUpdate; } /// @@ -68,7 +84,12 @@ public NitroSeedCoordinator( /// /// The name of the stage whose fusion configuration the gateways compose against. /// - public static NitroSeedCoordinator CreateProduction(string stage) + /// + /// Whether newly observed configurations are applied automatically by default. + /// + public static NitroSeedCoordinator CreateProduction( + string stage, + bool initialAutoUpdate = true) { ArgumentException.ThrowIfNullOrWhiteSpace(stage); @@ -94,13 +115,17 @@ public static NitroSeedCoordinator CreateProduction(string stage) GraphQLHttpClient.Create(httpClient, disposeHttpClient: false), timeProvider, NullLogger.Instance); + var stageUpdateClient = new NitroStageUpdateClient( + GraphQLHttpClient.Create(httpClient, disposeHttpClient: false)); return new NitroSeedCoordinator( stage, connectionResolver, seedProvider, schemaValidator, - NitroDefaults.CreateRunSeedDirectoryPath()); + stageUpdateClient, + NitroDefaults.CreateRunSeedDirectoryPath(), + initialAutoUpdate); } public Task ResolveConnectionAsync( @@ -125,6 +150,32 @@ public async Task ValidateSchemaAsync( cancellationToken); } + public async Task SubscribeToStageAsync( + string apiId, + ILogger logger, + CancellationToken cancellationToken) + { + var connection = await _connectionResolver.ResolveAsync(logger, cancellationToken); + return await _stageUpdateClient.SubscribeAsync( + connection, + apiId, + Stage, + cancellationToken); + } + + public async Task GetLatestStageSnapshotAsync( + string apiId, + ILogger logger, + CancellationToken cancellationToken) + { + var connection = await _connectionResolver.ResolveAsync(logger, cancellationToken); + return await _stageUpdateClient.GetLatestSnapshotAsync( + connection, + apiId, + Stage, + cancellationToken); + } + /// /// Acquires the fusion configuration of a gateway and keeps it as the configuration of that /// gateway for the rest of the run. @@ -165,16 +216,29 @@ public async Task AcquireSeedAsync( return NitroSeedAcquisition.Failed(result.Message!); } + var filePath = CopyToRunDirectory(gatewayName, result.FilePath!); var seed = new NitroGatewaySeed( apiId, Stage, - CopyToRunDirectory(gatewayName, result.FilePath!), + filePath, result.DownloadedAt!.Value, - result.Outcome is NitroSeedOutcome.Downloaded); + result.Outcome is NitroSeedOutcome.Downloaded, + await ComputeSchemaHashAsync(filePath, cancellationToken)); lock (_sync) { - _seedsByGateway[gatewayName] = seed; + if (_statesByGateway.TryGetValue(gatewayName, out var state)) + { + state.Current = seed; + state.Generation++; + state.Staged = null; + } + else + { + _statesByGateway.Add( + gatewayName, + new GatewaySeedState(seed, _initialAutoUpdate)); + } } return NitroSeedAcquisition.Acquired(seed); @@ -196,10 +260,211 @@ public async Task AcquireSeedAsync( lock (_sync) { - return _seedsByGateway.GetValueOrDefault(gatewayName); + return _statesByGateway.GetValueOrDefault(gatewayName)?.Current; + } + } + + public NitroSeedSnapshot? GetSeedSnapshot(string gatewayName) + { + ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); + + lock (_sync) + { + return _statesByGateway.TryGetValue(gatewayName, out var state) + ? new NitroSeedSnapshot(state.Current, state.Generation) + : null; + } + } + + public async Task DownloadFreshSeedAsync( + string gatewayName, + string apiId, + string versionIdentity, + ILogger logger, + bool suppressProviderLogs, + CancellationToken cancellationToken) + { + ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); + ArgumentException.ThrowIfNullOrWhiteSpace(apiId); + ArgumentException.ThrowIfNullOrWhiteSpace(versionIdentity); + ArgumentNullException.ThrowIfNull(logger); + + var connection = await _connectionResolver.ResolveAsync(logger, cancellationToken); + var providerLogger = suppressProviderLogs ? NullLogger.Instance : logger; + var result = await _seedProvider.GetSeedAsync( + connection, + apiId, + Stage, + providerLogger, + cancellationToken); + + if (result.Outcome is not NitroSeedOutcome.Downloaded) + { + return NitroSeedRefreshResult.Failed( + result.Message ?? "Nitro did not return a fresh Fusion configuration."); + } + + var filePath = CopyToRunDirectory(gatewayName, result.FilePath!); + var seed = new NitroGatewaySeed( + apiId, + Stage, + filePath, + result.DownloadedAt!.Value, + IsFresh: true, + await ComputeSchemaHashAsync(filePath, cancellationToken)); + + return NitroSeedRefreshResult.Downloaded( + new NitroSeedCandidate(seed, versionIdentity)); + } + + public bool IsAutoUpdateEnabled(string gatewayName) + { + ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); + + lock (_sync) + { + return _statesByGateway.TryGetValue(gatewayName, out var state) + ? state.AutoUpdate + : _initialAutoUpdate; + } + } + + public void SetAutoUpdate(string gatewayName, bool enabled) + { + ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); + + lock (_sync) + { + if (_statesByGateway.TryGetValue(gatewayName, out var state)) + { + state.AutoUpdate = enabled; + } } } + public void StageCandidate(string gatewayName, NitroSeedCandidate candidate) + { + ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); + ArgumentNullException.ThrowIfNull(candidate); + + lock (_sync) + { + if (_statesByGateway.TryGetValue(gatewayName, out var state)) + { + state.Staged = candidate; + } + } + } + + public NitroSeedAdoption? TryAdoptCandidate( + string gatewayName, + NitroSeedCandidate candidate, + bool wasStaged = false) + { + ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); + ArgumentNullException.ThrowIfNull(candidate); + + lock (_sync) + { + if (!_statesByGateway.TryGetValue(gatewayName, out var state)) + { + return null; + } + + var previous = new NitroSeedSnapshot(state.Current, state.Generation); + state.Current = candidate.Seed; + state.Generation++; + state.Staged = null; + + return new NitroSeedAdoption( + previous, + new NitroSeedSnapshot(state.Current, state.Generation), + candidate, + wasStaged); + } + } + + public NitroSeedAdoption? TryAdoptStaged(string gatewayName) + { + ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); + + lock (_sync) + { + if (!_statesByGateway.TryGetValue(gatewayName, out var state) + || state.Staged is not { } staged) + { + return null; + } + + var previous = new NitroSeedSnapshot(state.Current, state.Generation); + state.Current = staged.Seed; + state.Generation++; + state.Staged = null; + + return new NitroSeedAdoption( + previous, + new NitroSeedSnapshot(state.Current, state.Generation), + staged, + WasStaged: true); + } + } + + public void RollBackAdoption( + string gatewayName, + NitroSeedAdoption adoption, + bool restoreStaged) + { + ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); + ArgumentNullException.ThrowIfNull(adoption); + + lock (_sync) + { + if (!_statesByGateway.TryGetValue(gatewayName, out var state) + || state.Generation != adoption.Current.Generation) + { + return; + } + + state.Current = adoption.Previous.Seed; + state.Generation++; + state.Staged = restoreStaged ? adoption.Candidate : null; + } + } + + public NitroSeedCandidate? GetStagedCandidate(string gatewayName) + { + ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); + + lock (_sync) + { + return _statesByGateway.GetValueOrDefault(gatewayName)?.Staged; + } + } + + public NitroSeedCandidate? DiscardStagedCandidate(string gatewayName) + { + ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); + + lock (_sync) + { + if (!_statesByGateway.TryGetValue(gatewayName, out var state)) + { + return null; + } + + var staged = state.Staged; + state.Staged = null; + return staged; + } + } + + public void DeleteCandidate(NitroSeedCandidate candidate) + { + ArgumentNullException.ThrowIfNull(candidate); + + TryDelete(candidate.Seed.FilePath); + } + /// /// Deletes the private copies of this run. /// @@ -207,7 +472,7 @@ public void DeleteRunSeeds() { lock (_sync) { - _seedsByGateway.Clear(); + _statesByGateway.Clear(); } try @@ -227,10 +492,62 @@ private string CopyToRunDirectory(string gatewayName, string seedFilePath) { Directory.CreateDirectory(_runSeedDirectory); - var filePath = IOPath.Combine(_runSeedDirectory, gatewayName + ".far"); + var filePath = IOPath.Combine( + _runSeedDirectory, + $"{gatewayName}.{Interlocked.Increment(ref _nextRunSeedId):D8}.far"); File.Copy(seedFilePath, filePath, overwrite: true); return filePath; } + + private static async Task ComputeSchemaHashAsync( + string archivePath, + CancellationToken cancellationToken) + { + try + { + using var archive = FusionArchive.Open(archivePath); + using var configuration = await archive.TryGetGatewayConfigurationAsync( + WellKnownVersions.LatestGatewayFormatVersion, + cancellationToken); + + if (configuration is not null) + { + await using var schema = await configuration.OpenReadSchemaAsync(cancellationToken); + return Convert.ToHexString(await SHA256.HashDataAsync(schema, cancellationToken)); + } + } + catch (IOException) + { + // A seed produced by an older Nitro version can contain only source configurations. + // Hashing the complete immutable archive still gives the refresh path a stable guard. + } + + await using var archiveStream = File.OpenRead(archivePath); + return Convert.ToHexString(await SHA256.HashDataAsync(archiveStream, cancellationToken)); + } + + private static void TryDelete(string filePath) + { + try + { + File.Delete(filePath); + } + catch (Exception exception) when ( + exception is IOException or UnauthorizedAccessException) + { + } + } + + private sealed class GatewaySeedState(NitroGatewaySeed current, bool autoUpdate) + { + public NitroGatewaySeed Current { get; set; } = current; + + public NitroSeedCandidate? Staged { get; set; } + + public long Generation { get; set; } = 1; + + public bool AutoUpdate { get; set; } = autoUpdate; + } } diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateMonitor.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateMonitor.cs new file mode 100644 index 00000000000..8bcf590d1e6 --- /dev/null +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateMonitor.cs @@ -0,0 +1,481 @@ +using System.Threading.Channels; +using Microsoft.Extensions.Logging; +using Polly; +using Polly.Retry; + +namespace HotChocolate.Fusion.Aspire.Nitro; + +internal sealed class NitroSeedUpdateMonitor +{ + private static readonly TimeSpan s_initialReconnectDelay = TimeSpan.FromSeconds(10); + private static readonly TimeSpan s_maxReconnectDelay = TimeSpan.FromMinutes(5); + + private readonly Channel _updates = Channel.CreateBounded( + new BoundedChannelOptions(1) + { + FullMode = BoundedChannelFullMode.DropOldest, + SingleReader = true, + SingleWriter = true + }); + private readonly HashSet _reportedFailures = [with(StringComparer.Ordinal)]; + private readonly Lock _failureSync = new(); + private readonly string _gatewayName; + private readonly string _apiId; + private readonly NitroSeedCoordinator _coordinator; + private readonly SemaphoreSlim _compositionGate; + private readonly Func> _recomposeAsync; + private readonly ILogger _resourceLogger; + private readonly Action _notifyStaged; + private readonly Action _reportAdoption; + private readonly ResiliencePipeline _reconnectPipeline; + private readonly ILogger _logger; + private Task? _completion; + private bool _degraded; + private string? _lastProcessedIdentity; + + public NitroSeedUpdateMonitor( + string gatewayName, + string apiId, + NitroSeedCoordinator coordinator, + SemaphoreSlim compositionGate, + Func> recomposeAsync, + ILogger resourceLogger, + Action notifyStaged, + Action reportAdoption, + TimeProvider timeProvider, + ILogger logger) + { + ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); + ArgumentException.ThrowIfNullOrWhiteSpace(apiId); + ArgumentNullException.ThrowIfNull(coordinator); + ArgumentNullException.ThrowIfNull(compositionGate); + ArgumentNullException.ThrowIfNull(recomposeAsync); + ArgumentNullException.ThrowIfNull(resourceLogger); + ArgumentNullException.ThrowIfNull(notifyStaged); + ArgumentNullException.ThrowIfNull(reportAdoption); + ArgumentNullException.ThrowIfNull(timeProvider); + ArgumentNullException.ThrowIfNull(logger); + + _gatewayName = gatewayName; + _apiId = apiId; + _coordinator = coordinator; + _compositionGate = compositionGate; + _recomposeAsync = recomposeAsync; + _resourceLogger = resourceLogger; + _notifyStaged = notifyStaged; + _reportAdoption = reportAdoption; + _logger = logger; + + var pipelineBuilder = new ResiliencePipelineBuilder + { + TimeProvider = timeProvider + }; + _reconnectPipeline = pipelineBuilder + .AddRetry( + new RetryStrategyOptions + { + ShouldHandle = new PredicateBuilder().Handle(), + MaxRetryAttempts = int.MaxValue, + BackoffType = DelayBackoffType.Exponential, + Delay = s_initialReconnectDelay, + UseJitter = true, + MaxDelay = s_maxReconnectDelay, + DelayGenerator = arguments => new ValueTask( + CreateReconnectDelay(arguments.AttemptNumber)), + OnRetry = arguments => + { + RetryScheduled?.Invoke(arguments.RetryDelay); + return default; + } + }) + .Build(); + } + + internal Task Completion => _completion ?? Task.CompletedTask; + + internal Action? RetryScheduled { get; set; } + + public void Start(CancellationToken stoppingToken) + { + if (_completion is null) + { + _completion = RunAsync(stoppingToken); + } + } + + public async Task SetAutoUpdateAsync( + bool enabled, + CancellationToken cancellationToken) + { + await _compositionGate.WaitAsync(cancellationToken); + try + { + _coordinator.SetAutoUpdate(_gatewayName, enabled); + _resourceLogger.LogInformation( + "Automatic Nitro Fusion configuration updates were {AutoUpdateState} for " + + "{ResourceName}.", + enabled ? "enabled" : "disabled", + _gatewayName); + + if (!enabled || _coordinator.TryAdoptStaged(_gatewayName) is not { } adoption) + { + return; + } + + bool success; + try + { + success = await _recomposeAsync(adoption, cancellationToken); + } + catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) + { + _coordinator.RollBackAdoption( + _gatewayName, + adoption, + restoreStaged: true); + _lastProcessedIdentity = null; + throw; + } + catch (Exception exception) + { + _coordinator.RollBackAdoption( + _gatewayName, + adoption, + restoreStaged: true); + _lastProcessedIdentity = null; + _resourceLogger.LogWarning( + exception, + "The staged Nitro Fusion configuration for {ResourceName} could not be " + + "applied. The previous configuration remains active.", + _gatewayName); + return; + } + + if (!success) + { + _coordinator.RollBackAdoption( + _gatewayName, + adoption, + restoreStaged: true); + _lastProcessedIdentity = null; + return; + } + + _reportAdoption(adoption); + } + finally + { + _compositionGate.Release(); + } + } + + private async Task RunAsync(CancellationToken stoppingToken) + { + try + { + var producer = ProduceUpdatesAsync(stoppingToken); + var consumer = ConsumeUpdatesAsync(stoppingToken); + await Task.WhenAll(producer, consumer); + } + catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) + { + } + catch (Exception exception) + { + _logger.LogWarning( + exception, + "The Nitro stage update monitor for {ResourceName} stopped unexpectedly.", + _gatewayName); + } + } + + private async Task ProduceUpdatesAsync(CancellationToken stoppingToken) + { + var delayBeforeFirstAttempt = false; + + try + { + while (!stoppingToken.IsCancellationRequested) + { + var firstAttempt = true; + await _reconnectPipeline.ExecuteAsync( + async cancellationToken => + { + if (firstAttempt && delayBeforeFirstAttempt) + { + firstAttempt = false; + throw new NitroStageReconnectException(); + } + + firstAttempt = false; + if (!await RunConnectionCycleAsync(cancellationToken)) + { + throw new NitroStageReconnectException(); + } + }, + stoppingToken); + + // A connected SSE stream ended. Backend-initiated disconnects are normal, but + // the next connection still waits for the initial reconnect delay. Starting a + // new pipeline invocation resets the exponential failure count. + delayBeforeFirstAttempt = true; + } + } + finally + { + _updates.Writer.TryComplete(); + } + } + + private async Task RunConnectionCycleAsync(CancellationToken cancellationToken) + { + NitroStageSubscription? subscription = null; + try + { + subscription = await _coordinator.SubscribeToStageAsync( + _apiId, + _logger, + cancellationToken); + _degraded = false; + } + catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) + { + throw; + } + catch (Exception exception) + { + if (!_degraded) + { + _degraded = true; + _resourceLogger.LogWarning( + "The Nitro stage subscription for {ResourceName} is unavailable. Current " + + "stage versions will be checked during reconnect attempts.", + _gatewayName); + } + + _logger.LogDebug( + exception, + "The Nitro stage subscription for {ResourceName} could not be established.", + _gatewayName); + } + + NitroStageSnapshot? snapshot = null; + try + { + // The subscription is established before this query. Events remain buffered until + // enumeration starts, which closes the query-to-subscription publication race. + snapshot = await _coordinator.GetLatestStageSnapshotAsync( + _apiId, + _logger, + cancellationToken); + if (snapshot is not null) + { + _updates.Writer.TryWrite(snapshot); + } + } + catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) + { + throw; + } + catch (Exception exception) + { + ReportFailureOnce("query:" + NormalizeFailure(exception), exception); + } + + if (subscription is null) + { + return false; + } + + snapshot ??= new NitroStageSnapshot( + fusionConfigurationId: null, + new Dictionary>(StringComparer.Ordinal)); + + await using (subscription) + { + try + { + await foreach (var change in subscription.ReadChangesAsync(cancellationToken)) + { + snapshot = snapshot.Apply(change); + _updates.Writer.TryWrite(snapshot); + } + } + catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) + { + throw; + } + catch (Exception exception) + { + _logger.LogDebug( + exception, + "The Nitro stage subscription for {ResourceName} was interrupted.", + _gatewayName); + } + } + + return true; + } + + private async Task ConsumeUpdatesAsync(CancellationToken stoppingToken) + { + while (await _updates.Reader.WaitToReadAsync(stoppingToken)) + { + // Waiting for the composition gate before reading leaves the channel's one slot + // available to coalesce every stage change that arrives during a composition. + await _compositionGate.WaitAsync(stoppingToken); + try + { + if (!_updates.Reader.TryRead(out var update) + || string.Equals( + update.Identity, + _lastProcessedIdentity, + StringComparison.Ordinal)) + { + continue; + } + + if (await ProcessUpdateAsync(update, stoppingToken)) + { + _lastProcessedIdentity = update.Identity; + } + } + catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) + { + return; + } + catch (Exception exception) + { + ReportFailureOnce("process:" + NormalizeFailure(exception), exception); + } + finally + { + _compositionGate.Release(); + } + } + } + + private async Task ProcessUpdateAsync( + NitroStageSnapshot update, + CancellationToken cancellationToken) + { + var refresh = await _coordinator.DownloadFreshSeedAsync( + _gatewayName, + _apiId, + update.Identity, + _logger, + suppressProviderLogs: true, + cancellationToken); + + if (refresh.Candidate is not { } candidate) + { + ReportFailureOnce( + "download:" + NormalizeFailure(refresh.FailureMessage), + exception: null); + return false; + } + + if (_coordinator.GetSeedSnapshot(_gatewayName) is not { } current) + { + _coordinator.DeleteCandidate(candidate); + return false; + } + + if (string.Equals( + current.Seed.SchemaHash, + candidate.Seed.SchemaHash, + StringComparison.Ordinal)) + { + _coordinator.DeleteCandidate(candidate); + if (_coordinator.DiscardStagedCandidate(_gatewayName) is { } staged) + { + _coordinator.DeleteCandidate(staged); + } + + return true; + } + + if (!_coordinator.IsAutoUpdateEnabled(_gatewayName)) + { + _coordinator.StageCandidate(_gatewayName, candidate); + _notifyStaged(candidate); + return true; + } + + if (_coordinator.TryAdoptCandidate(_gatewayName, candidate) is not { } adoption) + { + return false; + } + + bool success; + try + { + success = await _recomposeAsync(adoption, cancellationToken); + } + catch + { + _coordinator.RollBackAdoption( + _gatewayName, + adoption, + restoreStaged: false); + throw; + } + + if (!success) + { + _coordinator.RollBackAdoption( + _gatewayName, + adoption, + restoreStaged: false); + return false; + } + + _reportAdoption(adoption); + return true; + } + + private void ReportFailureOnce(string reason, Exception? exception) + { + lock (_failureSync) + { + if (_reportedFailures.Count >= 20 || !_reportedFailures.Add(reason)) + { + return; + } + } + + _logger.LogWarning( + exception, + "Nitro stage update processing for {ResourceName} is unavailable: {Reason}", + _gatewayName, + reason); + } + + private static string NormalizeFailure(Exception exception) + => NormalizeFailure(exception.Message); + + private static string NormalizeFailure(string? message) + { + if (string.IsNullOrWhiteSpace(message)) + { + return "unknown"; + } + + var firstLine = message.Split(['\r', '\n'], 2)[0]; + return firstLine.Length <= 160 ? firstLine : firstLine[..160]; + } + + private static TimeSpan CreateReconnectDelay(int attemptNumber) + { + var exponent = Math.Min(attemptNumber, 20); + var exponentialMilliseconds = s_initialReconnectDelay.TotalMilliseconds + * Math.Pow(2, exponent); + var jitteredMilliseconds = exponentialMilliseconds + * (1 + Random.Shared.NextDouble() * 0.25); + + return TimeSpan.FromMilliseconds( + Math.Min(jitteredMilliseconds, s_maxReconnectDelay.TotalMilliseconds)); + } + + private sealed class NitroStageReconnectException : Exception; +} diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateNotifier.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateNotifier.cs new file mode 100644 index 00000000000..c4b7021e875 --- /dev/null +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateNotifier.cs @@ -0,0 +1,54 @@ +using Aspire.Hosting; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; + +namespace HotChocolate.Fusion.Aspire.Nitro; + +internal interface INitroSeedUpdateNotifier +{ + void NotifyAdopted(string message); + + void NotifyStaged(string message); +} + +#pragma warning disable ASPIREINTERACTION001 +internal sealed class NitroSeedUpdateNotifier( + IInteractionService interactionService, + IHostApplicationLifetime lifetime, + ILogger logger) + : INitroSeedUpdateNotifier +{ + public void NotifyAdopted(string message) => Notify(message); + + public void NotifyStaged(string message) => Notify(message); + + private void Notify(string message) + { + if (!interactionService.IsAvailable) + { + return; + } + + _ = PromptNotificationAsync(message); + } + + private async Task PromptNotificationAsync(string message) + { + try + { + await interactionService.PromptNotificationAsync( + "Nitro Fusion configuration", + message, + new NotificationInteractionOptions { Intent = MessageIntent.Information }, + lifetime.ApplicationStopping); + } + catch (OperationCanceledException) when (lifetime.ApplicationStopping.IsCancellationRequested) + { + } + catch (Exception exception) + { + logger.LogDebug(exception, "The Nitro Fusion configuration notification could not be shown."); + } + } +} +#pragma warning restore ASPIREINTERACTION001 diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateOptions.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateOptions.cs new file mode 100644 index 00000000000..68484d383b9 --- /dev/null +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateOptions.cs @@ -0,0 +1,21 @@ +namespace HotChocolate.Fusion.Aspire.Nitro; + +/// +/// Configures how Fusion Aspire follows changes to a Nitro stage during an AppHost run. +/// +public sealed class NitroSeedUpdateOptions +{ + /// + /// Gets or sets whether Fusion Aspire subscribes to stage changes and downloads newer Fusion + /// configurations. Enabling this option adds background subscription, query, and download + /// traffic to the configured Nitro API. The default is . + /// + public bool Enabled { get; set; } = true; + + /// + /// Gets or sets whether a newly downloaded Fusion configuration is immediately applied. When + /// disabled, new configurations are staged and applied by the next local recomposition. The + /// default is . + /// + public bool AutoUpdate { get; set; } = true; +} diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateService.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateService.cs new file mode 100644 index 00000000000..f22a83dcda0 --- /dev/null +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateService.cs @@ -0,0 +1,201 @@ +using Aspire.Hosting.ApplicationModel; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; + +namespace HotChocolate.Fusion.Aspire.Nitro; + +internal sealed class NitroSeedUpdateService +{ + private readonly Dictionary _monitors = [with(StringComparer.Ordinal)]; + private readonly Dictionary _resources = [with(StringComparer.Ordinal)]; + private readonly Dictionary> _notifiedVersions = [with(StringComparer.Ordinal)]; + private readonly HashSet _notifiedAdoptedHashes = [with(StringComparer.Ordinal)]; + private readonly Lock _sync = new(); + private readonly NitroCompositionOptions _options; + private readonly ResourceLoggerService _resourceLoggerService; + private readonly INitroSeedUpdateNotifier _notifier; + private readonly IHostApplicationLifetime _lifetime; + private readonly ILoggerFactory _loggerFactory; + private readonly TimeProvider _timeProvider; + + public NitroSeedUpdateService( + NitroCompositionOptions options, + ResourceLoggerService resourceLoggerService, + INitroSeedUpdateNotifier notifier, + IHostApplicationLifetime lifetime, + ILoggerFactory loggerFactory, + TimeProvider timeProvider) + { + _options = options; + _resourceLoggerService = resourceLoggerService; + _notifier = notifier; + _lifetime = lifetime; + _loggerFactory = loggerFactory; + _timeProvider = timeProvider; + } + + public bool IsEnabled => _options.Coordinator is not null && _options.SeedUpdates.Enabled; + + internal int MonitorCount + { + get + { + lock (_sync) + { + return _monitors.Count; + } + } + } + + public bool IsAutoUpdateEnabled(string gatewayName) + { + ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); + + return _options.Coordinator?.IsAutoUpdateEnabled(gatewayName) + ?? _options.SeedUpdates.AutoUpdate; + } + + public void Start( + IResource gateway, + string apiId, + SemaphoreSlim compositionGate, + Func> recomposeAsync) + { + ArgumentNullException.ThrowIfNull(gateway); + ArgumentException.ThrowIfNullOrWhiteSpace(apiId); + ArgumentNullException.ThrowIfNull(compositionGate); + ArgumentNullException.ThrowIfNull(recomposeAsync); + + if (!IsEnabled || _options.Coordinator is not { } coordinator) + { + return; + } + + NitroSeedUpdateMonitor monitor; + lock (_sync) + { + if (_monitors.ContainsKey(gateway.Name)) + { + return; + } + + monitor = new NitroSeedUpdateMonitor( + gateway.Name, + apiId, + coordinator, + compositionGate, + recomposeAsync, + _resourceLoggerService.GetLogger(gateway), + candidate => NotifyStaged(gateway.Name, coordinator.Stage, candidate), + adoption => ReportAdoption(gateway.Name, coordinator.Stage, adoption), + _timeProvider, + _loggerFactory.CreateLogger()); + _monitors.Add(gateway.Name, monitor); + _resources[gateway.Name] = gateway; + } + + // Starting outside the service lock prevents a synchronously completing test transport, + // or any future in-memory transport, from re-entering the reporting callbacks under it. + monitor.Start(_lifetime.ApplicationStopping); + } + + public async Task SetAutoUpdateAsync( + string gatewayName, + bool enabled, + CancellationToken cancellationToken) + { + NitroSeedUpdateMonitor? monitor; + lock (_sync) + { + _monitors.TryGetValue(gatewayName, out monitor); + } + + if (monitor is null) + { + return CommandResults.Failure("Nitro stage update monitoring is not ready."); + } + + await monitor.SetAutoUpdateAsync(enabled, cancellationToken); + + return CommandResults.Success( + enabled ? "Automatic Nitro updates enabled" : "Automatic Nitro updates disabled"); + } + + public void ReportAdoption( + string gatewayName, + string stage, + NitroSeedAdoption adoption) + { + ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); + ArgumentException.ThrowIfNullOrWhiteSpace(stage); + ArgumentNullException.ThrowIfNull(adoption); + + IResource? resource; + lock (_sync) + { + _resources.TryGetValue(gatewayName, out resource); + } + + if (resource is not null) + { + _resourceLoggerService.GetLogger(resource).LogInformation( + "Adopted a newer Fusion configuration for {ResourceName}: schema hash " + + "{PreviousHash} to {NewHash}, downloaded at {DownloadedAt}.", + gatewayName, + Prefix(adoption.Previous.Seed.SchemaHash), + Prefix(adoption.Current.Seed.SchemaHash), + adoption.Current.Seed.DownloadedAt.ToUniversalTime().ToString("yyyy-MM-dd HH:mm:ss") + "Z"); + } + + var shouldNotify = false; + lock (_sync) + { + var adoptedHashKey = gatewayName + "\n" + adoption.Current.Seed.SchemaHash; + if (_notifiedAdoptedHashes.Add(adoptedHashKey) + && MarkVersionNotified(gatewayName, adoption.Candidate.VersionIdentity)) + { + shouldNotify = true; + } + } + + if (shouldNotify) + { + _notifier.NotifyAdopted( + $"Recomposed '{gatewayName}' against a newer Fusion configuration " + + $"(stage '{stage}')."); + } + } + + private void NotifyStaged( + string gatewayName, + string stage, + NitroSeedCandidate candidate) + { + var shouldNotify = false; + lock (_sync) + { + shouldNotify = MarkVersionNotified(gatewayName, candidate.VersionIdentity); + } + + if (shouldNotify) + { + _notifier.NotifyStaged( + $"A newer Fusion configuration for '{gatewayName}' (stage '{stage}') was " + + "downloaded and staged. It is applied on the next recomposition."); + } + } + + private bool MarkVersionNotified(string gatewayName, string versionIdentity) + { + if (!_notifiedVersions.TryGetValue(gatewayName, out var identities)) + { + identities = new HashSet(StringComparer.Ordinal); + _notifiedVersions.Add(gatewayName, identities); + } + + return identities.Add(versionIdentity); + } + + private static string Prefix(string hash) + => hash.Length <= 12 ? hash : hash[..12]; +} diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateState.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateState.cs new file mode 100644 index 00000000000..1096f759371 --- /dev/null +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateState.cs @@ -0,0 +1,34 @@ +namespace HotChocolate.Fusion.Aspire.Nitro; + +internal sealed record NitroSeedCandidate( + NitroGatewaySeed Seed, + string VersionIdentity); + +internal sealed record NitroSeedSnapshot( + NitroGatewaySeed Seed, + long Generation); + +internal sealed record NitroSeedAdoption( + NitroSeedSnapshot Previous, + NitroSeedSnapshot Current, + NitroSeedCandidate Candidate, + bool WasStaged); + +internal sealed record NitroSeedRefreshResult( + NitroSeedCandidate? Candidate, + string? FailureMessage) +{ + public static NitroSeedRefreshResult Downloaded(NitroSeedCandidate candidate) + { + ArgumentNullException.ThrowIfNull(candidate); + + return new NitroSeedRefreshResult(candidate, FailureMessage: null); + } + + public static NitroSeedRefreshResult Failed(string failureMessage) + { + ArgumentException.ThrowIfNullOrWhiteSpace(failureMessage); + + return new NitroSeedRefreshResult(Candidate: null, failureMessage); + } +} diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroStageSnapshot.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroStageSnapshot.cs new file mode 100644 index 00000000000..f6a08ab57b4 --- /dev/null +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroStageSnapshot.cs @@ -0,0 +1,116 @@ +using System.Security.Cryptography; +using System.Text; + +namespace HotChocolate.Fusion.Aspire.Nitro; + +internal enum NitroStageChangeKind +{ + FusionConfigurationPublished, + ClientVersionPublished, + ClientVersionUnpublished, + ClientDeleted +} + +internal sealed record NitroStageChange( + NitroStageChangeKind Kind, + string? FusionConfigurationId = null, + string? ClientId = null, + string? ClientVersionId = null); + +internal sealed class NitroStageSnapshot +{ + private readonly Dictionary> _clientVersions; + + public NitroStageSnapshot( + string? fusionConfigurationId, + Dictionary> clientVersions) + { + FusionConfigurationId = fusionConfigurationId; + _clientVersions = clientVersions; + Identity = CreateIdentity(fusionConfigurationId, clientVersions); + } + + public string? FusionConfigurationId { get; } + + public string Identity { get; } + + public NitroStageSnapshot Apply(NitroStageChange change) + { + ArgumentNullException.ThrowIfNull(change); + + var fusionConfigurationId = FusionConfigurationId; + var clientVersions = CloneClientVersions(); + + switch (change.Kind) + { + case NitroStageChangeKind.FusionConfigurationPublished: + fusionConfigurationId = change.FusionConfigurationId; + break; + + case NitroStageChangeKind.ClientVersionPublished: + if (change.ClientId is { } publishedClientId + && change.ClientVersionId is { } publishedVersionId) + { + if (!clientVersions.TryGetValue(publishedClientId, out var versions)) + { + versions = new HashSet(StringComparer.Ordinal); + clientVersions.Add(publishedClientId, versions); + } + + versions.Add(publishedVersionId); + } + break; + + case NitroStageChangeKind.ClientVersionUnpublished: + if (change.ClientId is { } unpublishedClientId + && change.ClientVersionId is { } unpublishedVersionId + && clientVersions.TryGetValue(unpublishedClientId, out var publishedVersions)) + { + publishedVersions.Remove(unpublishedVersionId); + if (publishedVersions.Count == 0) + { + clientVersions.Remove(unpublishedClientId); + } + } + break; + + case NitroStageChangeKind.ClientDeleted: + if (change.ClientId is { } deletedClientId) + { + clientVersions.Remove(deletedClientId); + } + break; + } + + return new NitroStageSnapshot(fusionConfigurationId, clientVersions); + } + + private Dictionary> CloneClientVersions() + => _clientVersions.ToDictionary( + pair => pair.Key, + pair => new HashSet(pair.Value, StringComparer.Ordinal), + StringComparer.Ordinal); + + private static string CreateIdentity( + string? fusionConfigurationId, + Dictionary> clientVersions) + { + var builder = new StringBuilder(); + builder.Append("fusion:").Append(fusionConfigurationId).Append('\n'); + + foreach (var (clientId, versions) in clientVersions.OrderBy( + pair => pair.Key, + StringComparer.Ordinal)) + { + builder.Append("client:").Append(clientId).Append(':'); + foreach (var versionId in versions.Order(StringComparer.Ordinal)) + { + builder.Append(versionId).Append(','); + } + + builder.Append('\n'); + } + + return Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(builder.ToString()))); + } +} diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroStageUpdateClient.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroStageUpdateClient.cs new file mode 100644 index 00000000000..54884cfe9ba --- /dev/null +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroStageUpdateClient.cs @@ -0,0 +1,272 @@ +using System.Net.Http.Headers; +using System.Text.Json; +using HotChocolate.Language; +using HotChocolate.Transport; +using HotChocolate.Transport.Http; + +namespace HotChocolate.Fusion.Aspire.Nitro; + +internal sealed class NitroStageUpdateClient(GraphQLHttpClient client) + : INitroStageUpdateClient +{ + private const string EventStreamMediaType = "text/event-stream"; + + public async Task SubscribeAsync( + NitroConnection connection, + string apiId, + string stage, + CancellationToken cancellationToken) + { + ArgumentNullException.ThrowIfNull(connection); + ArgumentException.ThrowIfNullOrWhiteSpace(apiId); + ArgumentException.ThrowIfNullOrWhiteSpace(stage); + + // Fusion configuration publishes change the composition base. Client publish, unpublish, + // and delete events change the registered operations carried by the downloaded archive. + // The operation deliberately filters out OpenAPI and MCP changes, which do not affect a + // Fusion gateway archive. + var request = CreateRequest( + connection, + NitroOperationDocuments.WatchStageOperationName, +#if NITRO_PERSISTED_OPERATIONS + NitroOperationDocuments.GetWatchStageOperationId(), +#else + NitroOperationDocuments.GetWatchStageDocument(), +#endif + apiId, + stage); + request.Accept = GraphQLHttpRequest.GraphQLOverSse; + request.OperationKind = OperationType.Subscription; + + var response = await client.SendAsync(request, cancellationToken); + if (!response.IsSuccessStatusCode + || !string.Equals( + response.ContentHeaders.ContentType?.MediaType, + EventStreamMediaType, + StringComparison.OrdinalIgnoreCase)) + { + try + { + response.EnsureSuccessStatusCode(); + throw new InvalidOperationException( + "Nitro did not establish the stage-change SSE subscription."); + } + finally + { + response.Dispose(); + } + } + + return new HttpNitroStageSubscription(response); + } + + public async Task GetLatestSnapshotAsync( + NitroConnection connection, + string apiId, + string stage, + CancellationToken cancellationToken) + { + ArgumentNullException.ThrowIfNull(connection); + ArgumentException.ThrowIfNullOrWhiteSpace(apiId); + ArgumentException.ThrowIfNullOrWhiteSpace(stage); + + var request = CreateRequest( + connection, + NitroOperationDocuments.GetStageVersionOperationName, +#if NITRO_PERSISTED_OPERATIONS + NitroOperationDocuments.GetStageVersionOperationId(), +#else + NitroOperationDocuments.GetStageVersionDocument(), +#endif + apiId, + stage); + + using var response = await client.SendAsync(request, cancellationToken); + response.EnsureSuccessStatusCode(); + + using var result = await response.ReadAsResultAsync(cancellationToken); + ThrowOnErrors(result); + + if (result.Data.ValueKind is not JsonValueKind.Object + || !result.Data.TryGetProperty("apiById", out var api) + || api.ValueKind is not JsonValueKind.Object + || !api.TryGetProperty("stage", out var stageElement) + || stageElement.ValueKind is JsonValueKind.Null) + { + return null; + } + + if (stageElement.ValueKind is not JsonValueKind.Object) + { + throw new InvalidDataException("Nitro returned a malformed stage version."); + } + + var fusionConfigurationId = ReadNestedId( + stageElement, + "publishedFusionConfiguration"); + var clientVersions = new Dictionary>(StringComparer.Ordinal); + + if (stageElement.TryGetProperty("publishedClients", out var clients) + && clients.ValueKind is JsonValueKind.Array) + { + foreach (var publishedClient in clients.EnumerateArray()) + { + var clientId = ReadNestedId(publishedClient, "client"); + if (clientId is null) + { + continue; + } + + var versions = new HashSet(StringComparer.Ordinal); + if (publishedClient.TryGetProperty("publishedVersions", out var publishedVersions) + && publishedVersions.ValueKind is JsonValueKind.Array) + { + foreach (var publishedVersion in publishedVersions.EnumerateArray()) + { + if (ReadNestedId(publishedVersion, "version") is { } versionId) + { + versions.Add(versionId); + } + } + } + + if (versions.Count > 0) + { + clientVersions[clientId] = versions; + } + } + } + + return new NitroStageSnapshot(fusionConfigurationId, clientVersions); + } + + private static GraphQLHttpRequest CreateRequest( + NitroConnection connection, + string operationName, + string documentOrId, + string apiId, + string stage) + { + var variables = new Dictionary + { + ["apiId"] = apiId, + ["stageName"] = stage + }; + +#if NITRO_PERSISTED_OPERATIONS + var body = new OperationRequest( + id: documentOrId, + operationName: operationName, + variables: variables); +#else + var body = new OperationRequest( + documentOrId, + operationName: operationName, + variables: variables); +#endif + return new GraphQLHttpRequest(body, connection.GraphQLEndpoint) + { + OnMessageCreated = (_, requestMessage, _) => + NitroRequestHeaders.Apply(requestMessage, connection.Credential) + }; + } + + private static void ThrowOnErrors(OperationResult result) + { + if (result.Errors.ValueKind is JsonValueKind.Array + && result.Errors.GetArrayLength() > 0) + { + throw new InvalidDataException("Nitro returned GraphQL errors for a stage update operation."); + } + } + + private static string? ReadNestedId(JsonElement parent, string propertyName) + => parent.TryGetProperty(propertyName, out var value) + && value.ValueKind is JsonValueKind.Object + && value.TryGetProperty("id", out var id) + && id.ValueKind is JsonValueKind.String + ? id.GetString() + : null; + + private sealed class HttpNitroStageSubscription(GraphQLHttpResponse response) + : NitroStageSubscription + { + public override async IAsyncEnumerable ReadChangesAsync( + [System.Runtime.CompilerServices.EnumeratorCancellation] + CancellationToken cancellationToken) + { + await foreach (var result in response.ReadAsResultStreamAsync() + .WithCancellation(cancellationToken)) + { + using (result) + { + ThrowOnErrors(result); + + if (TryParseChange(result.Data) is { } change) + { + yield return change; + } + } + } + } + + public override ValueTask DisposeAsync() + { + response.Dispose(); + return ValueTask.CompletedTask; + } + + private static NitroStageChange? TryParseChange(JsonElement data) + { + if (data.ValueKind is not JsonValueKind.Object + || !data.TryGetProperty("onStageChanged", out var eventElement) + || eventElement.ValueKind is not JsonValueKind.Object + || !eventElement.TryGetProperty("__typename", out var typeNameElement) + || typeNameElement.ValueKind is not JsonValueKind.String) + { + throw new InvalidDataException("Nitro returned a malformed stage-change event."); + } + + return typeNameElement.GetString() switch + { + "FusionConfigurationPublishedStageChangeEvent" => new NitroStageChange( + NitroStageChangeKind.FusionConfigurationPublished, + FusionConfigurationId: ReadNestedId( + eventElement, + "fusionConfiguration")), + "ClientVersionPublishedStageChangeEvent" => CreateClientVersionChange( + eventElement, + NitroStageChangeKind.ClientVersionPublished), + "ClientVersionUnpublishedStageChangeEvent" => CreateClientVersionChange( + eventElement, + NitroStageChangeKind.ClientVersionUnpublished), + "ClientDeletedStageChangeEvent" => new NitroStageChange( + NitroStageChangeKind.ClientDeleted, + ClientId: ReadString(eventElement, "clientId")), + _ => null + }; + } + + private static NitroStageChange CreateClientVersionChange( + JsonElement eventElement, + NitroStageChangeKind kind) + { + if (!eventElement.TryGetProperty("clientVersion", out var clientVersion) + || clientVersion.ValueKind is not JsonValueKind.Object) + { + return new NitroStageChange(kind); + } + + return new NitroStageChange( + kind, + ClientId: ReadNestedId(clientVersion, "client"), + ClientVersionId: ReadString(clientVersion, "id")); + } + + private static string? ReadString(JsonElement parent, string propertyName) + => parent.TryGetProperty(propertyName, out var value) + && value.ValueKind is JsonValueKind.String + ? value.GetString() + : null; + } +} diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/Operations/GetNitroStageVersion.graphql b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/Operations/GetNitroStageVersion.graphql new file mode 100644 index 00000000000..c6853eb15eb --- /dev/null +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/Operations/GetNitroStageVersion.graphql @@ -0,0 +1,19 @@ +query GetNitroStageVersion($apiId: ID!, $stageName: String!) { + apiById(id: $apiId) { + stage(name: $stageName) { + publishedFusionConfiguration { + id + } + publishedClients { + client { + id + } + publishedVersions { + version { + id + } + } + } + } + } +} diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/Operations/GetNitroStageVersion.graphql.sha256 b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/Operations/GetNitroStageVersion.graphql.sha256 new file mode 100644 index 00000000000..e697fb4b8ba --- /dev/null +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/Operations/GetNitroStageVersion.graphql.sha256 @@ -0,0 +1 @@ +32cc7d1ee75aaa16627d21ee9289b4758f805d8b1eedfa796d832bae05b840f1 diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/Operations/WatchNitroStage.graphql b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/Operations/WatchNitroStage.graphql new file mode 100644 index 00000000000..cae141296bb --- /dev/null +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/Operations/WatchNitroStage.graphql @@ -0,0 +1,34 @@ +subscription WatchNitroStage($apiId: ID!, $stageName: String!) { + onStageChanged( + apiId: $apiId + stageName: $stageName + kind: [FUSION_CONFIGURATION, CLIENT] + ) { + __typename + kind + ... on FusionConfigurationPublishedStageChangeEvent { + fusionConfiguration { + id + } + } + ... on ClientVersionPublishedStageChangeEvent { + clientVersion { + id + client { + id + } + } + } + ... on ClientVersionUnpublishedStageChangeEvent { + clientVersion { + id + client { + id + } + } + } + ... on ClientDeletedStageChangeEvent { + clientId + } + } +} diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/Operations/WatchNitroStage.graphql.sha256 b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/Operations/WatchNitroStage.graphql.sha256 new file mode 100644 index 00000000000..0ac61da6ebe --- /dev/null +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/Operations/WatchNitroStage.graphql.sha256 @@ -0,0 +1 @@ +6626c4ef1d416a1db6a56bed1f19ab68f8797bb3e9bf79b31d2434b41a0d1d6e diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/NitroExtensions.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/NitroExtensions.cs index fbd72a5977c..4144e2972b4 100644 --- a/src/HotChocolate/Fusion/src/Fusion.Aspire/NitroExtensions.cs +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/NitroExtensions.cs @@ -1,6 +1,7 @@ using Aspire.Hosting; using Aspire.Hosting.ApplicationModel; using HotChocolate.Fusion.Aspire.Nitro; +using Microsoft.Extensions.DependencyInjection; namespace HotChocolate.Fusion.Aspire; @@ -64,6 +65,43 @@ public static IDistributedApplicationBuilder AddNitro( this IDistributedApplicationBuilder builder, string stage, Uri? portalUrl = null) + => AddNitroCore(builder, stage, portalUrl, configureSeedUpdates: null); + + /// + /// Adds GraphQL schema composition orchestration that composes against Nitro and configures + /// how Fusion Aspire follows changes to the selected stage during the AppHost run. + /// + /// The distributed application builder. + /// The Nitro stage whose Fusion configuration is used. + /// + /// Configures background stage-change subscriptions, current-version queries, Fusion + /// configuration downloads, and automatic adoption. + /// + /// + /// An optional Nitro portal URL. When omitted, the URL is derived from the effective Nitro API + /// URL. + /// + /// The distributed application builder for chaining. + /// + /// Stage update detection receives stage-change metadata and downloads the same Fusion archive + /// that startup seed acquisition downloads. It sends no schema or configuration data to Nitro. + /// + public static IDistributedApplicationBuilder AddNitro( + this IDistributedApplicationBuilder builder, + string stage, + Action configureSeedUpdates, + Uri? portalUrl = null) + { + ArgumentNullException.ThrowIfNull(configureSeedUpdates); + + return AddNitroCore(builder, stage, portalUrl, configureSeedUpdates); + } + + private static IDistributedApplicationBuilder AddNitroCore( + IDistributedApplicationBuilder builder, + string stage, + Uri? portalUrl, + Action? configureSeedUpdates) { ArgumentNullException.ThrowIfNull(builder); ArgumentException.ThrowIfNullOrWhiteSpace(stage); @@ -86,6 +124,7 @@ public static IDistributedApplicationBuilder AddNitro( } var options = SchemaCompositionRegistration.Ensure(builder); + configureSeedUpdates?.Invoke(options.SeedUpdates); if (options.Coordinator is { } coordinator) { @@ -106,11 +145,15 @@ public static IDistributedApplicationBuilder AddNitro( } options.PortalUrl ??= portalUrl; + AddAutoUpdateCommandsToConfiguredGateways(builder); return builder; } - options.Coordinator = NitroSeedCoordinator.CreateProduction(stage); + options.Coordinator = NitroSeedCoordinator.CreateProduction( + stage, + options.SeedUpdates.AutoUpdate); options.PortalUrl = portalUrl; + AddAutoUpdateCommandsToConfiguredGateways(builder); return builder; } @@ -147,6 +190,7 @@ public static IResourceBuilder WithNitroApiId( builder.WithAnnotation( new NitroApiIdAnnotation { ApiId = apiId }, ResourceAnnotationMutationBehavior.Replace); + TryAddAutoUpdateCommands(builder); return builder; } @@ -206,4 +250,88 @@ public static IResourceBuilder WithNitroSchemaValidation( internal static bool HasNitroSchemaValidation(this IResource resource) => resource.Annotations.OfType().Any(); + + internal static void TryAddAutoUpdateCommands(IResourceBuilder builder) + where T : IResource + { + if (!builder.Resource.NeedsGraphQLSchemaComposition() + || builder.Resource.GetNitroApiId() is null + || SchemaCompositionRegistration.GetOptions(builder.ApplicationBuilder)?.Coordinator + is null) + { + return; + } + + AddAutoUpdateCommand( + builder, + "disable-nitro-auto-update", + "Disable auto-update", + enabled: false); + AddAutoUpdateCommand( + builder, + "enable-nitro-auto-update", + "Enable auto-update", + enabled: true); + } + + private static void AddAutoUpdateCommandsToConfiguredGateways( + IDistributedApplicationBuilder builder) + { + foreach (var resource in builder.Resources.OfType()) + { + if (resource.NeedsGraphQLSchemaComposition() + && resource.GetNitroApiId() is not null) + { + TryAddAutoUpdateCommands(builder.CreateResourceBuilder(resource)); + } + } + } + + private static void AddAutoUpdateCommand( + IResourceBuilder builder, + string name, + string displayName, + bool enabled) + where T : IResource + { + if (builder.Resource.Annotations + .OfType() + .Any(command => command.Name == name)) + { + return; + } + + var resourceName = builder.Resource.Name; + + builder.WithCommand( + name, + displayName, + context => context.ServiceProvider + .GetService()? + .SetAutoUpdateAsync( + context.ResourceName, + enabled, + context.CancellationToken) + ?? Task.FromResult( + CommandResults.Failure("Nitro stage update monitoring is not ready.")), + new CommandOptions + { + Description = enabled + ? "Apply staged Nitro updates and resume automatic updates." + : "Stage Nitro updates until the next local recomposition.", + IconName = enabled ? "ArrowSyncCheckmark" : "ArrowSyncOff", + UpdateState = context => + { + var service = context.ServiceProvider.GetService(); + if (service is null || !service.IsEnabled) + { + return ResourceCommandState.Hidden; + } + + return service.IsAutoUpdateEnabled(resourceName) != enabled + ? ResourceCommandState.Enabled + : ResourceCommandState.Hidden; + } + }); + } } diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/SchemaComposition.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/SchemaComposition.cs index 73dc1cc7194..5d20ad04192 100644 --- a/src/HotChocolate/Fusion/src/Fusion.Aspire/SchemaComposition.cs +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/SchemaComposition.cs @@ -19,6 +19,7 @@ internal sealed class SchemaComposition( IHostApplicationLifetime lifetime, NitroCompositionOptions nitroOptions, NitroSchemaValidationCoordinator validationCoordinator, + NitroSeedUpdateService seedUpdateService, GatewayCompositionCommandCoordinator commandCoordinator, ILogger logger) : IDistributedApplicationEventingSubscriber @@ -289,6 +290,11 @@ internal async Task ComposeOnGatewayStartAsync( { validationCoordinator.Schedule(compositionResource, outcome.GatewaySchema); } + + StartSeedUpdateMonitor( + compositionResource, + appModel, + compositionGate); } /// @@ -510,6 +516,7 @@ internal async Task RunGuardedRecompositionAsync( outcome = await RecomposeSchemaAsync( compositionResource, appModel, + downloadFreshSeed: false, cancellationToken); } finally @@ -550,6 +557,7 @@ internal async Task ExecuteRecomposeCommandAsync( outcome = await RecomposeSchemaAsync( compositionResource, appModel, + downloadFreshSeed: true, cancellationToken); } catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) @@ -621,17 +629,30 @@ internal static Dictionary> private async Task RecomposeSchemaAsync( IResourceWithEndpoints compositionResource, DistributedApplicationModel appModel, + bool downloadFreshSeed, CancellationToken cancellationToken) { var coordinator = nitroOptions.Coordinator; var apiId = compositionResource.GetNitroApiId(); NitroGatewaySeed? seed = null; + NitroSeedAdoption? adoption = null; if (coordinator is not null && apiId is not null) { - // The fusion configuration of a run is fetched once, while the gateway starts, and - // every recomposition of the run builds on that same configuration. - seed = coordinator.GetSeed(compositionResource.Name); + if (downloadFreshSeed) + { + adoption = await TryPrepareManualSeedAdoptionAsync( + coordinator, + compositionResource, + apiId, + cancellationToken); + } + else + { + adoption = coordinator.TryAdoptStaged(compositionResource.Name); + } + + seed = adoption?.Current.Seed ?? coordinator.GetSeed(compositionResource.Name); if (seed is null) { @@ -647,11 +668,46 @@ private async Task RecomposeSchemaAsync( "Recomposing GraphQL schema for {ResourceName}...", compositionResource.Name); - var outcome = await ComposeSchemaAsync( - compositionResource, - appModel, - seed, - cancellationToken); + SchemaCompositionOutcome outcome; + try + { + outcome = await ComposeSchemaAsync( + compositionResource, + appModel, + seed, + cancellationToken); + } + catch + { + if (adoption is not null) + { + coordinator!.RollBackAdoption( + compositionResource.Name, + adoption, + restoreStaged: adoption.WasStaged); + } + + throw; + } + + if (adoption is not null) + { + if (outcome.Success) + { + seedUpdateService.ReportAdoption( + compositionResource.Name, + coordinator!.Stage, + adoption); + } + else + { + coordinator!.RollBackAdoption( + compositionResource.Name, + adoption, + restoreStaged: adoption.WasStaged); + } + } + if (outcome.Success) { logger.LogInformation( @@ -668,6 +724,133 @@ private async Task RecomposeSchemaAsync( return outcome; } + private async Task TryPrepareManualSeedAdoptionAsync( + NitroSeedCoordinator coordinator, + IResourceWithEndpoints compositionResource, + string apiId, + CancellationToken cancellationToken) + { + var gatewayLogger = CreateGatewayLogger(compositionResource); + + try + { + using var deadline = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + deadline.CancelAfter(TimeSpan.FromSeconds(30)); + + var refresh = await coordinator.DownloadFreshSeedAsync( + compositionResource.Name, + apiId, + versionIdentity: "manual", + gatewayLogger, + suppressProviderLogs: false, + deadline.Token); + + if (refresh.Candidate is { } downloaded + && coordinator.GetSeedSnapshot(compositionResource.Name) is { } current) + { + var candidate = downloaded with + { + VersionIdentity = "schema:" + downloaded.Seed.SchemaHash + }; + + if (string.Equals( + current.Seed.SchemaHash, + candidate.Seed.SchemaHash, + StringComparison.Ordinal)) + { + coordinator.DeleteCandidate(candidate); + if (coordinator.DiscardStagedCandidate(compositionResource.Name) is { } staged) + { + coordinator.DeleteCandidate(staged); + } + + return null; + } + + return coordinator.TryAdoptCandidate(compositionResource.Name, candidate); + } + + gatewayLogger.LogWarning( + "A fresh Fusion configuration could not be downloaded before manually " + + "recomposing {ResourceName}. The held configuration will be used.", + compositionResource.Name); + } + catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) + { + throw; + } + catch (Exception exception) + { + gatewayLogger.LogWarning( + exception, + "A fresh Fusion configuration could not be downloaded before manually " + + "recomposing {ResourceName}. The held configuration will be used.", + compositionResource.Name); + } + + return coordinator.TryAdoptStaged(compositionResource.Name); + } + + private void StartSeedUpdateMonitor( + IResourceWithEndpoints compositionResource, + DistributedApplicationModel appModel, + SemaphoreSlim compositionGate) + { + if (compositionResource.GetNitroApiId() is not { } apiId) + { + return; + } + + seedUpdateService.Start( + compositionResource, + apiId, + compositionGate, + (adoption, cancellationToken) => RecomposeAdoptedSeedAsync( + compositionResource, + appModel, + adoption, + cancellationToken)); + } + + private async Task RecomposeAdoptedSeedAsync( + IResourceWithEndpoints compositionResource, + DistributedApplicationModel appModel, + NitroSeedAdoption adoption, + CancellationToken cancellationToken) + { + logger.LogInformation( + "Recomposing GraphQL schema for {ResourceName} against an updated Nitro " + + "configuration...", + compositionResource.Name); + + var outcome = await ComposeSchemaAsync( + compositionResource, + appModel, + adoption.Current.Seed, + cancellationToken); + + if (!outcome.Success) + { + logger.LogWarning( + "Schema recomposition for {ResourceName} against the updated Nitro " + + "configuration failed. The gateway keeps the previous schema.", + compositionResource.Name); + return false; + } + + logger.LogInformation( + "Schema recomposition for {ResourceName} against the updated Nitro configuration " + + "completed.", + compositionResource.Name); + + if (outcome.GatewaySchema is not null) + { + validationCoordinator.Schedule(compositionResource, outcome.GatewaySchema); + } + + return true; + } + private async Task ComposeSchemaAsync( IResourceWithEndpoints compositionResource, DistributedApplicationModel appModel, diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/SchemaCompositionRegistration.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/SchemaCompositionRegistration.cs index 41854507ec0..85481be68a2 100644 --- a/src/HotChocolate/Fusion/src/Fusion.Aspire/SchemaCompositionRegistration.cs +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/SchemaCompositionRegistration.cs @@ -33,7 +33,10 @@ public static NitroCompositionOptions Ensure(IDistributedApplicationBuilder buil builder.Services.TryAddEventingSubscriber(); builder.Services.TryAddSingleton(); builder.Services.TryAddSingleton(); + builder.Services.TryAddSingleton(); + builder.Services.TryAddSingleton(); builder.Services.TryAddSingleton(); + builder.Services.TryAddSingleton(TimeProvider.System); return options; } diff --git a/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/CompositionHarness.cs b/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/CompositionHarness.cs index 16f24a2419e..d46f81de57e 100644 --- a/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/CompositionHarness.cs +++ b/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/CompositionHarness.cs @@ -34,18 +34,27 @@ public static CompositionHarness Create( Coordinator = coordinator, PortalUrl = portalUrl }; + options.SeedUpdates.Enabled = false; var validationCoordinator = new NitroSchemaValidationCoordinator( options, resourceLoggerService, notifier, lifetime, NullLoggerFactory.Instance); + var seedUpdateService = new NitroSeedUpdateService( + options, + resourceLoggerService, + NoopSeedUpdateNotifier.Instance, + lifetime, + NullLoggerFactory.Instance, + TimeProvider.System); var composition = new SchemaComposition( notifications, resourceLoggerService, lifetime, options, validationCoordinator, + seedUpdateService, new GatewayCompositionCommandCoordinator(), logger); @@ -57,6 +66,19 @@ public static CompositionHarness Create( lifetime); } + private sealed class NoopSeedUpdateNotifier : INitroSeedUpdateNotifier + { + public static NoopSeedUpdateNotifier Instance { get; } = new(); + + public void NotifyAdopted(string message) + { + } + + public void NotifyStaged(string message) + { + } + } + private sealed class EmptyServiceProvider : IServiceProvider { public static EmptyServiceProvider Instance { get; } = new(); diff --git a/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroOperationDocumentsTests.cs b/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroOperationDocumentsTests.cs index bc24cf1982e..cebf4f9fae7 100644 --- a/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroOperationDocumentsTests.cs +++ b/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroOperationDocumentsTests.cs @@ -52,6 +52,22 @@ public void ValidationOperationName_Should_MatchTheOperationInTheDocument( Assert.Contains(expectedDeclaration, document, StringComparison.Ordinal); } + [Theory] + [InlineData("GetNitroStageVersion", "query GetNitroStageVersion(")] + [InlineData("WatchNitroStage", "subscription WatchNitroStage(")] + public void StageUpdateOperationName_Should_MatchTheOperationInTheDocument( + string operationName, + string expectedDeclaration) + { + // act + var document = operationName == NitroOperationDocuments.GetStageVersionOperationName + ? NitroOperationDocuments.GetStageVersionDocument() + : NitroOperationDocuments.GetWatchStageDocument(); + + // assert + Assert.Contains(expectedDeclaration, document, StringComparison.Ordinal); + } + #endif #if NITRO_PERSISTED_OPERATIONS @@ -86,5 +102,25 @@ public void GetValidationOperationId_Should_ReturnTheEmbeddedHash( // assert Assert.Equal(expectedOperationId, operationId); } + + [Theory] + [InlineData( + "GetNitroStageVersion", + "32cc7d1ee75aaa16627d21ee9289b4758f805d8b1eedfa796d832bae05b840f1")] + [InlineData( + "WatchNitroStage", + "6626c4ef1d416a1db6a56bed1f19ab68f8797bb3e9bf79b31d2434b41a0d1d6e")] + public void GetStageUpdateOperationId_Should_ReturnTheEmbeddedHash( + string operationName, + string expectedOperationId) + { + // act + var operationId = operationName == NitroOperationDocuments.GetStageVersionOperationName + ? NitroOperationDocuments.GetStageVersionOperationId() + : NitroOperationDocuments.GetWatchStageOperationId(); + + // assert + Assert.Equal(expectedOperationId, operationId); + } #endif } diff --git a/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroSchemaCompositionTests.cs b/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroSchemaCompositionTests.cs index 45dc2b0c289..093bb0ca597 100644 --- a/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroSchemaCompositionTests.cs +++ b/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroSchemaCompositionTests.cs @@ -720,6 +720,114 @@ await WaitForValidationReportAsync( """); } + [Fact] + public async Task ExecuteRecomposeCommandAsync_Should_DownloadAndAdoptFreshSeedBeforeComposition() + { + // arrange + await ServeSeedAsync(); + var coordinator = CreateCoordinator(); + var harness = CompositionHarness.Create(coordinator); + var (model, gateway) = CreateModel(GatewayApiId); + using var compositionGate = new SemaphoreSlim(1, 1); + await coordinator.AcquireSeedAsync( + gateway.Name, + GatewayApiId, + new RecordingLogger(), + TestContext.Current.CancellationToken); + var previousHash = coordinator.GetSeed(gateway.Name)!.SchemaHash; + await ServeUpdatedSeedAsync(); + + // act + var result = await harness.Composition.ExecuteRecomposeCommandAsync( + gateway, + model, + compositionGate, + TestContext.Current.CancellationToken); + + // assert + $""" + Result: {result.Success} ({result.Message}) + Downloads: {GetDownloadCount()} + Source schemas: {await ReadSourceSchemaNamesAsync()} + Seed changed: {coordinator.GetSeed(gateway.Name)!.SchemaHash != previousHash} + """.MatchInlineSnapshot( + """ + Result: True (Schema composition completed) + Downloads: 2 + Source schemas: inventory, products + Seed changed: True + """); + } + + [Fact] + public async Task ExecuteRecomposeCommandAsync_Should_FallBackToHeldSeed_WhenDownloadFails() + { + // arrange + await ServeSeedAsync(); + var coordinator = CreateCoordinator(); + var harness = CompositionHarness.Create(coordinator); + var (model, gateway) = CreateModel(GatewayApiId); + using var compositionGate = new SemaphoreSlim(1, 1); + await coordinator.AcquireSeedAsync( + gateway.Name, + GatewayApiId, + new RecordingLogger(), + TestContext.Current.CancellationToken); + _server.DownloadHandler = _ => + FakeNitroResponse.Status(StatusCodes.Status503ServiceUnavailable); + + // act + var result = await harness.Composition.ExecuteRecomposeCommandAsync( + gateway, + model, + compositionGate, + TestContext.Current.CancellationToken); + + // assert + Assert.Equal( + "True|2|orders, products, reviews", + $"{result.Success}|{GetDownloadCount()}|{await ReadSourceSchemaNamesAsync()}"); + } + + [Fact] + public async Task RunGuardedRecompositionAsync_Should_AdoptStagedSeedWithoutDownloading() + { + // arrange + await ServeSeedAsync(); + var coordinator = CreateCoordinator(); + var harness = CompositionHarness.Create(coordinator); + var (model, gateway) = CreateModel(GatewayApiId); + using var compositionGate = new SemaphoreSlim(1, 1); + await coordinator.AcquireSeedAsync( + gateway.Name, + GatewayApiId, + new RecordingLogger(), + TestContext.Current.CancellationToken); + await ServeUpdatedSeedAsync(); + var refresh = await coordinator.DownloadFreshSeedAsync( + gateway.Name, + GatewayApiId, + "configuration-2", + new RecordingLogger(), + suppressProviderLogs: false, + TestContext.Current.CancellationToken); + coordinator.StageCandidate(gateway.Name, refresh.Candidate!); + var downloadsBeforeComposition = GetDownloadCount(); + + // act + await harness.Composition.RunGuardedRecompositionAsync( + gateway, + model, + compositionGate, + TestContext.Current.CancellationToken); + + // assert + Assert.Equal( + $"{downloadsBeforeComposition}|{downloadsBeforeComposition}|inventory, products|False", + $"{downloadsBeforeComposition}|{GetDownloadCount()}|{await ReadSourceSchemaNamesAsync()}|" + + $"{coordinator.GetStagedCandidate(gateway.Name) is not null}"); + } + private (DistributedApplicationModel Model, IResourceWithEndpoints Gateway) CreateModel( string? gatewayApiId, string? productsApiId = null, @@ -789,13 +897,33 @@ private NitroSeedCoordinator CreateCoordinator( GraphQLHttpClient.Create(_httpClient, disposeHttpClient: false), _timeProvider, new RecordingLogger()), - _directory.GetPath("run")); + new NoopStageUpdateClient(), + _directory.GetPath("run"), + initialAutoUpdate: true); private TestNitroEnvironment CreateDefaultEnvironment() => new( (NitroEnvironmentVariables.CloudUrl, _server.BaseAddress.AbsoluteUri), (NitroEnvironmentVariables.ApiKey, "nitro-api-key")); + private sealed class NoopStageUpdateClient : INitroStageUpdateClient + { + public Task SubscribeAsync( + NitroConnection connection, + string apiId, + string stage, + CancellationToken cancellationToken) + => Task.FromException( + new InvalidOperationException("The composition harness does not run a stage monitor.")); + + public Task GetLatestSnapshotAsync( + NitroConnection connection, + string apiId, + string stage, + CancellationToken cancellationToken) + => Task.FromResult(null); + } + private static int CountRecompositions(CompositionHarness harness) => harness.Logger.Entries.Count( entry => entry.Level is LogLevel.Information @@ -879,6 +1007,22 @@ private async Task ServeSeedAsync() _server.DownloadHandler = _ => FakeNitroResponse.Archive(archive); } + private async Task ServeUpdatedSeedAsync() + { + var archive = await NitroTestArchive.CreateAsync( + TestContext.Current.CancellationToken, + new NitroTestSourceSchema( + "products", + "type Query { staleProduct: String }", + CreateSettings("products", "https://stale.example.com/graphql")), + new NitroTestSourceSchema( + "inventory", + "type Query { inventory: Int }", + CreateSettings("inventory", "https://inventory.example.com/graphql"))); + + _server.DownloadHandler = _ => FakeNitroResponse.Archive(archive); + } + private async Task PrimeTheCacheAsync() { await ServeSeedAsync(); diff --git a/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroSeedUpdateMonitorTests.cs b/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroSeedUpdateMonitorTests.cs new file mode 100644 index 00000000000..019aac4f8d3 --- /dev/null +++ b/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroSeedUpdateMonitorTests.cs @@ -0,0 +1,556 @@ +using System.Collections.Concurrent; +using System.Net; +using HotChocolate.Transport.Http; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Time.Testing; + +namespace HotChocolate.Fusion.Aspire.Nitro; + +public sealed class NitroSeedUpdateMonitorTests : IAsyncLifetime +{ + private const string ApiId = "QXBpCmdhdGV3YXk"; + private const string Stage = "production"; + private readonly NitroTestDirectory _directory = new(); + private byte[] _initialArchive = null!; + private byte[] _updatedArchive = null!; + + public async ValueTask InitializeAsync() + { + _initialArchive = await CreateArchiveAsync("type Query { value: String }"); + _updatedArchive = await CreateArchiveAsync("type Query { value: Int }"); + } + + public ValueTask DisposeAsync() + { + _directory.Dispose(); + return ValueTask.CompletedTask; + } + + [Fact] + public async Task RunAsync_Should_SubscribeBeforeQueryAndDeduplicateQueryEventVersion() + { + // arrange + var stageClient = new ScriptedStageUpdateClient(); + var queried = Snapshot("configuration-2"); + stageClient.EnqueueSnapshot(queried); + stageClient.EnqueueSubscription( + new NitroStageChange( + NitroStageChangeKind.FusionConfigurationPublished, + FusionConfigurationId: "configuration-2")); + var handler = new ArchiveSequenceHandler(_initialArchive, _updatedArchive); + var coordinator = await CreateCoordinatorAsync(handler, stageClient, autoUpdate: true); + using var gate = new SemaphoreSlim(1, 1); + using var stopping = new CancellationTokenSource(); + var recompositions = new List(); + var adoptions = new List(); + var monitor = CreateMonitor( + coordinator, + gate, + adoption => + { + recompositions.Add(adoption); + return Task.FromResult(true); + }, + staged => { }, + adoptions.Add, + TimeProvider.System); + + // act + monitor.Start(stopping.Token); + await WaitUntilAsync(() => adoptions.Count == 1); + await stopping.CancelAsync(); + await monitor.Completion.WaitAsync( + TimeSpan.FromSeconds(5), + TestContext.Current.CancellationToken); + + // assert + $""" + Operations: {string.Join(", ", stageClient.Operations.Take(3))} + Downloads: {handler.RequestCount} + Recompositions: {recompositions.Count} + Adoptions: {adoptions.Count} + Version: {adoptions[0].Candidate.VersionIdentity == queried.Identity} + """.MatchInlineSnapshot( + """ + Operations: subscribe, query, event:FusionConfigurationPublished + Downloads: 2 + Recompositions: 1 + Adoptions: 1 + Version: True + """); + } + + [Fact] + public async Task RunAsync_Should_SkipAdoption_WhenDownloadedHashIsUnchanged() + { + // arrange + var stageClient = new ScriptedStageUpdateClient(); + stageClient.EnqueueSnapshot(Snapshot("configuration-1")); + stageClient.EnqueueSubscription(); + var handler = new ArchiveSequenceHandler(_initialArchive, _initialArchive); + var coordinator = await CreateCoordinatorAsync(handler, stageClient, autoUpdate: true); + using var gate = new SemaphoreSlim(1, 1); + using var stopping = new CancellationTokenSource(); + var recompositions = 0; + var notifications = 0; + var monitor = CreateMonitor( + coordinator, + gate, + _ => + { + recompositions++; + return Task.FromResult(true); + }, + _ => notifications++, + _ => notifications++, + TimeProvider.System); + + // act + monitor.Start(stopping.Token); + await WaitUntilAsync(() => handler.RequestCount == 2); + await stopping.CancelAsync(); + await monitor.Completion.WaitAsync( + TimeSpan.FromSeconds(5), + TestContext.Current.CancellationToken); + + // assert + Assert.Equal("2|0|0", $"{handler.RequestCount}|{recompositions}|{notifications}"); + } + + [Fact] + public async Task RunAsync_Should_ProcessOnlyNewestVersion_WhenCompositionGateIsHeld() + { + // arrange + var stageClient = new ScriptedStageUpdateClient(); + var first = Snapshot("configuration-1"); + var secondChange = new NitroStageChange( + NitroStageChangeKind.FusionConfigurationPublished, + FusionConfigurationId: "configuration-2"); + var thirdChange = new NitroStageChange( + NitroStageChangeKind.FusionConfigurationPublished, + FusionConfigurationId: "configuration-3"); + var newest = first.Apply(secondChange).Apply(thirdChange); + stageClient.EnqueueSnapshot(first); + stageClient.EnqueueSubscription(secondChange, thirdChange); + var handler = new ArchiveSequenceHandler(_initialArchive, _updatedArchive); + var coordinator = await CreateCoordinatorAsync(handler, stageClient, autoUpdate: true); + using var gate = new SemaphoreSlim(1, 1); + await gate.WaitAsync(TestContext.Current.CancellationToken); + using var stopping = new CancellationTokenSource(); + var adoptions = new List(); + var monitor = CreateMonitor( + coordinator, + gate, + _ => Task.FromResult(true), + _ => { }, + adoptions.Add, + TimeProvider.System); + + // act + monitor.Start(stopping.Token); + await WaitUntilAsync(() => stageClient.Operations.Count >= 4); + gate.Release(); + await WaitUntilAsync(() => adoptions.Count == 1); + await stopping.CancelAsync(); + await monitor.Completion.WaitAsync( + TimeSpan.FromSeconds(5), + TestContext.Current.CancellationToken); + + // assert + Assert.Equal( + $"2|1|{newest.Identity}", + $"{handler.RequestCount}|{adoptions.Count}|{adoptions[0].Candidate.VersionIdentity}"); + } + + [Fact] + public async Task RunAsync_Should_QueryAndAdopt_WhenSubscriptionCannotBeEstablished() + { + // arrange + var stageClient = new ScriptedStageUpdateClient(); + stageClient.EnqueueSnapshot(Snapshot("configuration-2")); + stageClient.EnqueueSubscriptionFailure(new HttpRequestException("SSE unavailable")); + var handler = new ArchiveSequenceHandler(_initialArchive, _updatedArchive); + var coordinator = await CreateCoordinatorAsync(handler, stageClient, autoUpdate: true); + using var gate = new SemaphoreSlim(1, 1); + using var stopping = new CancellationTokenSource(); + var adoptions = 0; + var monitor = CreateMonitor( + coordinator, + gate, + _ => Task.FromResult(true), + _ => { }, + _ => adoptions++, + TimeProvider.System); + + // act + monitor.Start(stopping.Token); + await WaitUntilAsync(() => adoptions == 1); + await stopping.CancelAsync(); + await monitor.Completion.WaitAsync( + TimeSpan.FromSeconds(5), + TestContext.Current.CancellationToken); + + // assert + Assert.Equal("subscribe, query|2|1", $"{string.Join(", ", stageClient.Operations)}|{handler.RequestCount}|{adoptions}"); + } + + [Fact] + public async Task SetAutoUpdateAsync_Should_AdoptStagedVersion_WhenEnabled() + { + // arrange + var stageClient = new ScriptedStageUpdateClient(); + stageClient.EnqueueSnapshot(Snapshot("configuration-2")); + stageClient.EnqueueBlockingSubscription(); + var handler = new ArchiveSequenceHandler(_initialArchive, _updatedArchive); + var coordinator = await CreateCoordinatorAsync(handler, stageClient, autoUpdate: false); + using var gate = new SemaphoreSlim(1, 1); + using var stopping = new CancellationTokenSource(); + var staged = 0; + var recompositions = 0; + var adoptions = 0; + var monitor = CreateMonitor( + coordinator, + gate, + _ => + { + recompositions++; + return Task.FromResult(true); + }, + _ => staged++, + _ => adoptions++, + TimeProvider.System); + + // act + monitor.Start(stopping.Token); + await WaitUntilAsync(() => staged == 1); + var beforeEnable = coordinator.GetStagedCandidate("gateway") is not null; + await monitor.SetAutoUpdateAsync( + enabled: true, + TestContext.Current.CancellationToken); + await stopping.CancelAsync(); + await monitor.Completion.WaitAsync( + TimeSpan.FromSeconds(5), + TestContext.Current.CancellationToken); + + // assert + $""" + Staged before enable: {beforeEnable} + Staged notifications: {staged} + Recompositions: {recompositions} + Adoptions: {adoptions} + Auto-update: {coordinator.IsAutoUpdateEnabled("gateway")} + """.MatchInlineSnapshot( + """ + Staged before enable: True + Staged notifications: 1 + Recompositions: 1 + Adoptions: 1 + Auto-update: True + """); + } + + [Fact] + public async Task RunAsync_Should_GrowBackoffAndResetAfterSuccessfulConnection() + { + // arrange + var timeProvider = new FakeTimeProvider(); + var stageClient = new ScriptedStageUpdateClient(); + stageClient.EnqueueSnapshot(null); + stageClient.EnqueueSnapshot(null); + stageClient.EnqueueSnapshot(null); + stageClient.EnqueueSubscriptionFailure(new HttpRequestException("first")); + stageClient.EnqueueSubscriptionFailure(new HttpRequestException("second")); + stageClient.EnqueueSubscription(); + var handler = new ArchiveSequenceHandler(_initialArchive); + var coordinator = await CreateCoordinatorAsync(handler, stageClient, autoUpdate: true); + using var gate = new SemaphoreSlim(1, 1); + using var stopping = new CancellationTokenSource(); + var delays = new ConcurrentQueue(); + var monitor = CreateMonitor( + coordinator, + gate, + _ => Task.FromResult(true), + _ => { }, + _ => { }, + timeProvider); + monitor.RetryScheduled = delays.Enqueue; + + // act + monitor.Start(stopping.Token); + await WaitUntilAsync(() => delays.Count == 1); + timeProvider.Advance(TimeSpan.FromSeconds(13)); + await WaitUntilAsync(() => delays.Count == 2); + timeProvider.Advance(TimeSpan.FromSeconds(26)); + await WaitUntilAsync(() => delays.Count == 3); + await stopping.CancelAsync(); + await monitor.Completion.WaitAsync( + TimeSpan.FromSeconds(5), + TestContext.Current.CancellationToken); + + // assert + var observed = delays.ToArray(); + Assert.InRange(observed[0], TimeSpan.FromSeconds(10), TimeSpan.FromSeconds(12.5)); + Assert.InRange(observed[1], TimeSpan.FromSeconds(20), TimeSpan.FromSeconds(25)); + Assert.InRange(observed[2], TimeSpan.FromSeconds(10), TimeSpan.FromSeconds(12.5)); + } + + [Fact] + public async Task RunAsync_Should_RetryVersionAfterAdoptionFailsOnPreviousCycle() + { + // arrange + var timeProvider = new FakeTimeProvider(); + var stageClient = new ScriptedStageUpdateClient(); + var version = Snapshot("configuration-2"); + stageClient.EnqueueSnapshot(version); + stageClient.EnqueueSnapshot(version); + stageClient.EnqueueSubscription(); + stageClient.EnqueueBlockingSubscription(); + var handler = new ArchiveSequenceHandler( + _initialArchive, + _updatedArchive, + _updatedArchive); + var coordinator = await CreateCoordinatorAsync(handler, stageClient, autoUpdate: true); + var originalHash = coordinator.GetSeed("gateway")!.SchemaHash; + using var gate = new SemaphoreSlim(1, 1); + using var stopping = new CancellationTokenSource(); + var attempts = 0; + var adoptions = 0; + var monitor = CreateMonitor( + coordinator, + gate, + _ => Task.FromResult(++attempts > 1), + _ => { }, + _ => adoptions++, + timeProvider); + + // act + monitor.Start(stopping.Token); + await WaitUntilAsync(() => attempts == 1); + var rolledBack = coordinator.GetSeed("gateway")!.SchemaHash == originalHash; + timeProvider.Advance(TimeSpan.FromSeconds(13)); + await WaitUntilAsync(() => attempts == 2); + await stopping.CancelAsync(); + await monitor.Completion.WaitAsync( + TimeSpan.FromSeconds(5), + TestContext.Current.CancellationToken); + + // assert + Assert.Equal("True|3|2|1", $"{rolledBack}|{handler.RequestCount}|{attempts}|{adoptions}"); + } + + private async Task CreateCoordinatorAsync( + ArchiveSequenceHandler handler, + INitroStageUpdateClient stageClient, + bool autoUpdate) + { + var timeProvider = TimeProvider.System; + var httpClient = new HttpClient(handler); + var connectionResolver = new NitroConnectionResolver( + new NitroSessionReader(_directory.GetPath("session.json"), TimeSpan.Zero), + new TestNitroEnvironment( + (NitroEnvironmentVariables.CloudUrl, "https://nitro.example.test"), + (NitroEnvironmentVariables.ApiKey, "key")), + NitroDefaults.ApiUrl, + timeProvider, + NitroDefaults.AccessTokenExpiryGrace); + var coordinator = new NitroSeedCoordinator( + Stage, + connectionResolver, + new NitroSeedProvider( + new NitroFusionConfigurationDownloader( + httpClient, + new NitroDownloadRetryPolicy(1, 1, TimeSpan.Zero), + timeProvider), + new NitroSeedCache(_directory.GetPath("cache"), timeProvider), + new NitroApiLookupClient( + GraphQLHttpClient.Create(httpClient, disposeHttpClient: false))), + NoopSchemaValidator.Instance, + stageClient, + _directory.GetPath("run"), + autoUpdate); + + var acquisition = await coordinator.AcquireSeedAsync( + "gateway", + ApiId, + new RecordingLogger(), + TestContext.Current.CancellationToken); + if (acquisition.Seed is null) + { + throw new InvalidOperationException(acquisition.FailureMessage); + } + + return coordinator; + } + + private static NitroSeedUpdateMonitor CreateMonitor( + NitroSeedCoordinator coordinator, + SemaphoreSlim gate, + Func> recomposeAsync, + Action notifyStaged, + Action reportAdoption, + TimeProvider timeProvider) + => new( + "gateway", + ApiId, + coordinator, + gate, + (adoption, _) => recomposeAsync(adoption), + new RecordingLogger(), + notifyStaged, + reportAdoption, + timeProvider, + new RecordingLogger()); + + private static NitroStageSnapshot Snapshot(string? configurationId) + => new( + configurationId, + new Dictionary>(StringComparer.Ordinal)); + + private static Task CreateArchiveAsync(string schema) + => NitroTestArchive.CreateAsync( + TestContext.Current.CancellationToken, + new NitroTestSourceSchema( + "products", + schema, + """ + { + "name": "products", + "transports": { + "http": { + "url": "https://products.example.test/graphql" + } + } + } + """)); + + private static async Task WaitUntilAsync(Func condition) + { + using var timeout = new CancellationTokenSource(TimeSpan.FromSeconds(5)); + while (!condition()) + { + await Task.Delay(10, timeout.Token); + } + } + + private sealed class ArchiveSequenceHandler(params byte[][] archives) : HttpMessageHandler + { + private readonly ConcurrentQueue _archives = new(archives); + private byte[]? _lastArchive = archives.LastOrDefault(); + private int _requestCount; + + public int RequestCount => Volatile.Read(ref _requestCount); + + protected override Task SendAsync( + HttpRequestMessage request, + CancellationToken cancellationToken) + { + Interlocked.Increment(ref _requestCount); + if (_archives.TryDequeue(out var archive)) + { + _lastArchive = archive; + } + + return Task.FromResult( + new HttpResponseMessage(HttpStatusCode.OK) + { + Content = new ByteArrayContent(_lastArchive ?? []) + }); + } + } + + private sealed class ScriptedStageUpdateClient : INitroStageUpdateClient + { + private readonly ConcurrentQueue _subscriptions = new(); + private readonly ConcurrentQueue _snapshots = new(); + + public ConcurrentQueue Operations { get; } = new(); + + public void EnqueueSnapshot(NitroStageSnapshot? snapshot) + => _snapshots.Enqueue(snapshot); + + public void EnqueueSubscription(params NitroStageChange[] changes) + => _subscriptions.Enqueue(new SubscriptionPlan(changes, Block: false, Failure: null)); + + public void EnqueueBlockingSubscription() + => _subscriptions.Enqueue(new SubscriptionPlan([], Block: true, Failure: null)); + + public void EnqueueSubscriptionFailure(Exception exception) + => _subscriptions.Enqueue(new SubscriptionPlan([], Block: false, exception)); + + public Task SubscribeAsync( + NitroConnection connection, + string apiId, + string stage, + CancellationToken cancellationToken) + { + Operations.Enqueue("subscribe"); + if (!_subscriptions.TryDequeue(out var plan)) + { + plan = new SubscriptionPlan([], Block: true, Failure: null); + } + + return plan.Failure is null + ? Task.FromResult( + new ScriptedSubscription(plan, Operations)) + : Task.FromException(plan.Failure); + } + + public Task GetLatestSnapshotAsync( + NitroConnection connection, + string apiId, + string stage, + CancellationToken cancellationToken) + { + Operations.Enqueue("query"); + _snapshots.TryDequeue(out var snapshot); + return Task.FromResult(snapshot); + } + + private sealed class ScriptedSubscription( + SubscriptionPlan plan, + ConcurrentQueue operations) + : NitroStageSubscription + { + public override async IAsyncEnumerable ReadChangesAsync( + [System.Runtime.CompilerServices.EnumeratorCancellation] + CancellationToken cancellationToken) + { + foreach (var change in plan.Changes) + { + operations.Enqueue("event:" + change.Kind); + yield return change; + } + + if (plan.Block) + { + await Task.Delay(Timeout.InfiniteTimeSpan, cancellationToken); + } + } + + public override ValueTask DisposeAsync() => ValueTask.CompletedTask; + } + + private sealed record SubscriptionPlan( + IReadOnlyList Changes, + bool Block, + Exception? Failure); + } + + private sealed class NoopSchemaValidator : INitroSchemaValidator + { + public static NoopSchemaValidator Instance { get; } = new(); + + public Task ValidateAsync( + NitroConnection connection, + string apiId, + string stage, + byte[] schema, + string schemaHash, + CancellationToken cancellationToken) + => Task.FromResult( + NitroSchemaValidationReport.Passed( + schemaHash, + requestId: "noop", + DateTimeOffset.UtcNow)); + } +} diff --git a/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroStageUpdateClientTests.cs b/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroStageUpdateClientTests.cs new file mode 100644 index 00000000000..c11adca3603 --- /dev/null +++ b/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroStageUpdateClientTests.cs @@ -0,0 +1,144 @@ +using System.Collections.Concurrent; +using System.Net; +using System.Net.Http.Headers; +using System.Text; +using HotChocolate.Transport.Http; + +namespace HotChocolate.Fusion.Aspire.Nitro; + +public sealed class NitroStageUpdateClientTests +{ + [Fact] + public async Task SubscribeAndQuery_Should_PreserveEventPublishedDuringQueryRace() + { + // arrange + var handler = new ResponseSequenceHandler( + Sse( + """ + event: next + data: {"data":{"onStageChanged":{"__typename":"ClientVersionPublishedStageChangeEvent","kind":"CLIENT","clientVersion":{"id":"client-version-2","client":{"id":"client-1"}}}}} + + event: complete + + """), + Json( + """ + { + "data": { + "apiById": { + "stage": { + "publishedFusionConfiguration": { + "id": "configuration-1" + }, + "publishedClients": [ + { + "client": { + "id": "client-1" + }, + "publishedVersions": [ + { + "version": { + "id": "client-version-2" + } + } + ] + } + ] + } + } + } + } + """)); + using var httpClient = new HttpClient(handler); + var client = new NitroStageUpdateClient( + GraphQLHttpClient.Create(httpClient, disposeHttpClient: false)); + var connection = new NitroConnection( + new Uri("https://nitro.example.test"), + new Uri("https://nitro.example.test/graphql"), + NitroCredential.FromApiKey("secret")); + + // act + await using var subscription = await client.SubscribeAsync( + connection, + "api-1", + "production", + TestContext.Current.CancellationToken); + var snapshot = await client.GetLatestSnapshotAsync( + connection, + "api-1", + "production", + TestContext.Current.CancellationToken); + var changes = new List(); + await foreach (var change in subscription.ReadChangesAsync( + TestContext.Current.CancellationToken)) + { + changes.Add(change); + } + + // assert + var changedSnapshot = snapshot!.Apply(Assert.Single(changes)); + $""" + Requests: {handler.Requests.Count} + First operation: {ReadOperation(handler.Requests[0].Body)} + Second operation: {ReadOperation(handler.Requests[1].Body)} + API key sent: {handler.Requests.All(request => request.ApiKey == "secret")} + Identity unchanged: {changedSnapshot.Identity == snapshot.Identity} + """.MatchInlineSnapshot( + """ + Requests: 2 + First operation: WatchNitroStage + Second operation: GetNitroStageVersion + API key sent: True + Identity unchanged: True + """); + } + + private static string ReadOperation(string body) + { + using var document = System.Text.Json.JsonDocument.Parse(body); + return document.RootElement.GetProperty("operationName").GetString()!; + } + + private static HttpResponseMessage Json(string body) + => new(HttpStatusCode.OK) + { + Content = new StringContent(body, Encoding.UTF8, "application/json") + }; + + private static HttpResponseMessage Sse(string body) + { + var content = new StringContent(body, Encoding.UTF8); + content.Headers.ContentType = new MediaTypeHeaderValue("text/event-stream") + { + CharSet = "utf-8" + }; + + return new HttpResponseMessage(HttpStatusCode.OK) { Content = content }; + } + + private sealed class ResponseSequenceHandler(params HttpResponseMessage[] responses) + : HttpMessageHandler + { + private readonly ConcurrentQueue _responses = new(responses); + + public List Requests { get; } = []; + + protected override async Task SendAsync( + HttpRequestMessage request, + CancellationToken cancellationToken) + { + Requests.Add( + new RecordedGraphQLRequest( + await request.Content!.ReadAsStringAsync(cancellationToken), + request.Headers.TryGetValues(NitroRequestHeaders.ApiKey, out var values) + ? values.Single() + : null)); + + return _responses.TryDequeue(out var response) + ? response + : new HttpResponseMessage(HttpStatusCode.InternalServerError); + } + } + + private sealed record RecordedGraphQLRequest(string Body, string? ApiKey); +} diff --git a/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/NitroExtensionsTests.cs b/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/NitroExtensionsTests.cs index 198011feff4..d2f8941d752 100644 --- a/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/NitroExtensionsTests.cs +++ b/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/NitroExtensionsTests.cs @@ -4,6 +4,7 @@ using Aspire.Hosting.Lifecycle; using HotChocolate.Fusion.Aspire.Nitro; using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging.Abstractions; using IOPath = System.IO.Path; @@ -338,6 +339,154 @@ public void AddNitro_Should_StoreTheCallerSuppliedPortalUrl() Assert.Same(portalUrl, GetNitroCompositionOptions(builder).PortalUrl); } + [Fact] + public void AddNitro_Should_ConfigureSeedUpdates() + { + // arrange + var builder = DistributedApplication.CreateBuilder(); + + // act + builder.AddNitro( + "production", + options => + { + options.Enabled = false; + options.AutoUpdate = false; + }); + + // assert + var options = GetNitroCompositionOptions(builder).SeedUpdates; + Assert.Equal("False|False", $"{options.Enabled}|{options.AutoUpdate}"); + } + + [Fact] + public void NitroGateway_Should_RegisterBothAutoUpdateCommands_WhenNitroIsAddedFirst() + { + // arrange + var builder = DistributedApplication.CreateBuilder(); + builder.AddNitro("production"); + + // act + var gateway = builder + .AddProject("gateway", GetTestProjectFile()) + .WithGraphQLSchemaComposition() + .WithNitroApiId("QXBpCmdhdGV3YXk"); + + // assert + string.Join( + Environment.NewLine, + gateway.Resource.Annotations + .OfType() + .Select(command => $"{command.Name}: {command.DisplayName}") + .Order(StringComparer.Ordinal)) + .MatchInlineSnapshot( + """ + disable-nitro-auto-update: Disable auto-update + enable-nitro-auto-update: Enable auto-update + recompose: Recompose + """); + } + + [Fact] + public void NitroGateway_Should_RegisterBothAutoUpdateCommands_WhenNitroIsAddedLast() + { + // arrange + var builder = DistributedApplication.CreateBuilder(); + var gateway = builder + .AddProject("gateway", GetTestProjectFile()) + .WithGraphQLSchemaComposition() + .WithNitroApiId("QXBpCmdhdGV3YXk"); + + // act + builder.AddNitro("production"); + + // assert + Assert.Equal( + 2, + gateway.Resource.Annotations + .OfType() + .Count(command => command.Name.Contains("nitro-auto-update", StringComparison.Ordinal))); + } + + [Fact] + public void SeedUpdateService_Should_NotStartMonitor_WhenDetectionIsDisabled() + { + // arrange + var builder = DistributedApplication.CreateBuilder(); + builder.AddNitro( + "production", + options => options.Enabled = false); + var gateway = builder + .AddProject("gateway", GetTestProjectFile()) + .WithGraphQLSchemaComposition() + .WithNitroApiId("QXBpCmdhdGV3YXk"); + var lifetime = new TestHostApplicationLifetime(); + var resourceLoggerService = new ResourceLoggerService(); + var service = new NitroSeedUpdateService( + GetNitroCompositionOptions(builder), + resourceLoggerService, + NoopSeedUpdateNotifier.Instance, + lifetime, + NullLoggerFactory.Instance, + TimeProvider.System); + using var gate = new SemaphoreSlim(1, 1); + + // act + service.Start( + gateway.Resource, + "QXBpCmdhdGV3YXk", + gate, + (_, _) => Task.FromResult(true)); + + // assert + Assert.Equal(0, service.MonitorCount); + } + + [Fact] + public async Task AutoUpdateCommands_Should_ShowOnlyDisableCommand_WhenAutoUpdateStartsEnabled() + { + // arrange + var builder = DistributedApplication.CreateBuilder(); + builder.AddNitro("production"); + var gateway = builder + .AddProject("gateway", GetTestProjectFile()) + .WithGraphQLSchemaComposition() + .WithNitroApiId("QXBpCmdhdGV3YXk"); + var lifetime = new TestHostApplicationLifetime(); + var service = new NitroSeedUpdateService( + GetNitroCompositionOptions(builder), + new ResourceLoggerService(), + NoopSeedUpdateNotifier.Instance, + lifetime, + NullLoggerFactory.Instance, + TimeProvider.System); + await using var services = new ServiceCollection() + .AddSingleton(service) + .BuildServiceProvider(); + var commands = gateway.Resource.Annotations + .OfType() + .Where(command => command.Name.Contains("nitro-auto-update", StringComparison.Ordinal)) + .OrderBy(command => command.Name, StringComparer.Ordinal) + .ToArray(); + var context = new UpdateCommandStateContext + { + ResourceSnapshot = new CustomResourceSnapshot + { + ResourceType = "project", + Properties = [] + }, + ServiceProvider = services + }; + + // act + var states = commands.Select(command => command.UpdateState!(context)); + + // assert + Assert.Equal( + [ResourceCommandState.Enabled, ResourceCommandState.Hidden], + states); + } + /// /// Describes every registration of the schema composition. The distributed application /// registers eventing subscribers of its own, so only the registrations of the composition @@ -363,4 +512,17 @@ private static string GetTestProjectFile([CallerFilePath] string sourceFile = "" => IOPath.Combine( IOPath.GetDirectoryName(sourceFile)!, "HotChocolate.Fusion.Aspire.Tests.csproj"); + + private sealed class NoopSeedUpdateNotifier : INitroSeedUpdateNotifier + { + public static NoopSeedUpdateNotifier Instance { get; } = new(); + + public void NotifyAdopted(string message) + { + } + + public void NotifyStaged(string message) + { + } + } } From aea65cba7c7af432e4d750141c2dee8b0332f5ab Mon Sep 17 00:00:00 2001 From: Michael Staib Date: Fri, 31 Jul 2026 15:31:45 +0200 Subject: [PATCH 2/2] edits --- .../Nitro/NitroSeedCoordinator.cs | 78 ++++++++++++++++++- .../Nitro/NitroSeedUpdateMonitor.cs | 2 + .../Nitro/NitroSeedUpdateService.cs | 10 +++ .../src/Fusion.Aspire/NitroExtensions.cs | 20 +++-- .../src/Fusion.Aspire/SchemaComposition.cs | 5 +- .../Nitro/NitroSeedUpdateMonitorTests.cs | 8 +- .../NitroExtensionsTests.cs | 57 ++++++++++++-- 7 files changed, 161 insertions(+), 19 deletions(-) diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedCoordinator.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedCoordinator.cs index 6716b2af06f..906867451d5 100644 --- a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedCoordinator.cs +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedCoordinator.cs @@ -21,7 +21,7 @@ internal sealed class NitroSeedCoordinator private readonly INitroSchemaValidator _schemaValidator; private readonly INitroStageUpdateClient _stageUpdateClient; private readonly string _runSeedDirectory; - private readonly bool _initialAutoUpdate; + private bool _initialAutoUpdate; private long _nextRunSeedId; /// @@ -225,10 +225,14 @@ public async Task AcquireSeedAsync( result.Outcome is NitroSeedOutcome.Downloaded, await ComputeSchemaHashAsync(filePath, cancellationToken)); + NitroGatewaySeed? replacedSeed = null; + NitroSeedCandidate? discardedStaged = null; lock (_sync) { if (_statesByGateway.TryGetValue(gatewayName, out var state)) { + replacedSeed = state.Current; + discardedStaged = state.Staged; state.Current = seed; state.Generation++; state.Staged = null; @@ -241,6 +245,9 @@ public async Task AcquireSeedAsync( } } + TryDeleteUnlessCurrent(replacedSeed?.FilePath, seed.FilePath); + TryDeleteUnlessCurrent(discardedStaged?.Seed.FilePath, seed.FilePath); + return NitroSeedAcquisition.Acquired(seed); } @@ -329,6 +336,19 @@ public bool IsAutoUpdateEnabled(string gatewayName) } } + public void SetInitialAutoUpdate(bool enabled) + { + lock (_sync) + { + _initialAutoUpdate = enabled; + + foreach (var state in _statesByGateway.Values) + { + state.AutoUpdate = enabled; + } + } + } + public void SetAutoUpdate(string gatewayName, bool enabled) { ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); @@ -347,13 +367,17 @@ public void StageCandidate(string gatewayName, NitroSeedCandidate candidate) ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); ArgumentNullException.ThrowIfNull(candidate); + NitroSeedCandidate? replaced = null; lock (_sync) { if (_statesByGateway.TryGetValue(gatewayName, out var state)) { + replaced = state.Staged; state.Staged = candidate; } } + + TryDeleteUnlessCurrent(replaced?.Seed.FilePath, candidate.Seed.FilePath); } public NitroSeedAdoption? TryAdoptCandidate( @@ -364,6 +388,8 @@ public void StageCandidate(string gatewayName, NitroSeedCandidate candidate) ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); ArgumentNullException.ThrowIfNull(candidate); + NitroSeedCandidate? discardedStaged; + NitroSeedAdoption adoption; lock (_sync) { if (!_statesByGateway.TryGetValue(gatewayName, out var state)) @@ -371,17 +397,21 @@ public void StageCandidate(string gatewayName, NitroSeedCandidate candidate) return null; } + discardedStaged = state.Staged; var previous = new NitroSeedSnapshot(state.Current, state.Generation); state.Current = candidate.Seed; state.Generation++; state.Staged = null; - return new NitroSeedAdoption( + adoption = new NitroSeedAdoption( previous, new NitroSeedSnapshot(state.Current, state.Generation), candidate, wasStaged); } + + TryDeleteUnlessCurrent(discardedStaged?.Seed.FilePath, candidate.Seed.FilePath); + return adoption; } public NitroSeedAdoption? TryAdoptStaged(string gatewayName) @@ -417,6 +447,7 @@ public void RollBackAdoption( ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); ArgumentNullException.ThrowIfNull(adoption); + var deleteCandidate = false; lock (_sync) { if (!_statesByGateway.TryGetValue(gatewayName, out var state) @@ -428,6 +459,40 @@ public void RollBackAdoption( state.Current = adoption.Previous.Seed; state.Generation++; state.Staged = restoreStaged ? adoption.Candidate : null; + deleteCandidate = !restoreStaged; + } + + if (deleteCandidate) + { + TryDeleteUnlessCurrent( + adoption.Current.Seed.FilePath, + adoption.Previous.Seed.FilePath); + } + } + + public void CompleteAdoption( + string gatewayName, + NitroSeedAdoption adoption) + { + ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); + ArgumentNullException.ThrowIfNull(adoption); + + var deletePrevious = false; + lock (_sync) + { + deletePrevious = _statesByGateway.TryGetValue(gatewayName, out var state) + && state.Generation == adoption.Current.Generation + && string.Equals( + state.Current.FilePath, + adoption.Current.Seed.FilePath, + StringComparison.Ordinal); + } + + if (deletePrevious) + { + TryDeleteUnlessCurrent( + adoption.Previous.Seed.FilePath, + adoption.Current.Seed.FilePath); } } @@ -540,6 +605,15 @@ private static void TryDelete(string filePath) } } + private static void TryDeleteUnlessCurrent(string? filePath, string currentFilePath) + { + if (filePath is not null + && !string.Equals(filePath, currentFilePath, StringComparison.Ordinal)) + { + TryDelete(filePath); + } + } + private sealed class GatewaySeedState(NitroGatewaySeed current, bool autoUpdate) { public NitroGatewaySeed Current { get; set; } = current; diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateMonitor.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateMonitor.cs index 8bcf590d1e6..23a18025812 100644 --- a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateMonitor.cs +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateMonitor.cs @@ -161,6 +161,7 @@ public async Task SetAutoUpdateAsync( return; } + _coordinator.CompleteAdoption(_gatewayName, adoption); _reportAdoption(adoption); } finally @@ -430,6 +431,7 @@ private async Task ProcessUpdateAsync( return false; } + _coordinator.CompleteAdoption(_gatewayName, adoption); _reportAdoption(adoption); return true; } diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateService.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateService.cs index f22a83dcda0..a8c498fba51 100644 --- a/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateService.cs +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/Nitro/NitroSeedUpdateService.cs @@ -55,6 +55,16 @@ public bool IsAutoUpdateEnabled(string gatewayName) ?? _options.SeedUpdates.AutoUpdate; } + public bool IsReady(string gatewayName) + { + ArgumentException.ThrowIfNullOrWhiteSpace(gatewayName); + + lock (_sync) + { + return _monitors.ContainsKey(gatewayName); + } + } + public void Start( IResource gateway, string apiId, diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/NitroExtensions.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/NitroExtensions.cs index 4144e2972b4..d88995750cd 100644 --- a/src/HotChocolate/Fusion/src/Fusion.Aspire/NitroExtensions.cs +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/NitroExtensions.cs @@ -73,14 +73,14 @@ public static IDistributedApplicationBuilder AddNitro( /// /// The distributed application builder. /// The Nitro stage whose Fusion configuration is used. - /// - /// Configures background stage-change subscriptions, current-version queries, Fusion - /// configuration downloads, and automatic adoption. - /// /// /// An optional Nitro portal URL. When omitted, the URL is derived from the effective Nitro API /// URL. /// + /// + /// Configures background stage-change subscriptions, current-version queries, Fusion + /// configuration downloads, and automatic adoption. + /// /// The distributed application builder for chaining. /// /// Stage update detection receives stage-change metadata and downloads the same Fusion archive @@ -89,8 +89,8 @@ public static IDistributedApplicationBuilder AddNitro( public static IDistributedApplicationBuilder AddNitro( this IDistributedApplicationBuilder builder, string stage, - Action configureSeedUpdates, - Uri? portalUrl = null) + Uri? portalUrl, + Action configureSeedUpdates) { ArgumentNullException.ThrowIfNull(configureSeedUpdates); @@ -124,7 +124,6 @@ private static IDistributedApplicationBuilder AddNitroCore( } var options = SchemaCompositionRegistration.Ensure(builder); - configureSeedUpdates?.Invoke(options.SeedUpdates); if (options.Coordinator is { } coordinator) { @@ -144,11 +143,14 @@ private static IDistributedApplicationBuilder AddNitroCore( $"Nitro is already added with the portal URL '{options.PortalUrl}'."); } + configureSeedUpdates?.Invoke(options.SeedUpdates); + coordinator.SetInitialAutoUpdate(options.SeedUpdates.AutoUpdate); options.PortalUrl ??= portalUrl; AddAutoUpdateCommandsToConfiguredGateways(builder); return builder; } + configureSeedUpdates?.Invoke(options.SeedUpdates); options.Coordinator = NitroSeedCoordinator.CreateProduction( stage, options.SeedUpdates.AutoUpdate); @@ -323,7 +325,9 @@ private static void AddAutoUpdateCommand( UpdateState = context => { var service = context.ServiceProvider.GetService(); - if (service is null || !service.IsEnabled) + if (service is null + || !service.IsEnabled + || !service.IsReady(resourceName)) { return ResourceCommandState.Hidden; } diff --git a/src/HotChocolate/Fusion/src/Fusion.Aspire/SchemaComposition.cs b/src/HotChocolate/Fusion/src/Fusion.Aspire/SchemaComposition.cs index 5d20ad04192..cc84ff2229f 100644 --- a/src/HotChocolate/Fusion/src/Fusion.Aspire/SchemaComposition.cs +++ b/src/HotChocolate/Fusion/src/Fusion.Aspire/SchemaComposition.cs @@ -694,9 +694,12 @@ private async Task RecomposeSchemaAsync( { if (outcome.Success) { + coordinator!.CompleteAdoption( + compositionResource.Name, + adoption); seedUpdateService.ReportAdoption( compositionResource.Name, - coordinator!.Stage, + coordinator.Stage, adoption); } else diff --git a/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroSeedUpdateMonitorTests.cs b/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroSeedUpdateMonitorTests.cs index 019aac4f8d3..6531206d698 100644 --- a/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroSeedUpdateMonitorTests.cs +++ b/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/Nitro/NitroSeedUpdateMonitorTests.cs @@ -39,6 +39,7 @@ public async Task RunAsync_Should_SubscribeBeforeQueryAndDeduplicateQueryEventVe FusionConfigurationId: "configuration-2")); var handler = new ArchiveSequenceHandler(_initialArchive, _updatedArchive); var coordinator = await CreateCoordinatorAsync(handler, stageClient, autoUpdate: true); + var initialSeedPath = coordinator.GetSeed("gateway")!.FilePath; using var gate = new SemaphoreSlim(1, 1); using var stopping = new CancellationTokenSource(); var recompositions = new List(); @@ -70,6 +71,7 @@ await monitor.Completion.WaitAsync( Recompositions: {recompositions.Count} Adoptions: {adoptions.Count} Version: {adoptions[0].Candidate.VersionIdentity == queried.Identity} + Retained seeds: {File.Exists(initialSeedPath)}|{File.Exists(coordinator.GetSeed("gateway")!.FilePath)} """.MatchInlineSnapshot( """ Operations: subscribe, query, event:FusionConfigurationPublished @@ -77,6 +79,7 @@ await monitor.Completion.WaitAsync( Recompositions: 1 Adoptions: 1 Version: True + Retained seeds: False|True """); } @@ -327,6 +330,7 @@ public async Task RunAsync_Should_RetryVersionAfterAdoptionFailsOnPreviousCycle( monitor.Start(stopping.Token); await WaitUntilAsync(() => attempts == 1); var rolledBack = coordinator.GetSeed("gateway")!.SchemaHash == originalHash; + var filesAfterRollback = Directory.GetFiles(_directory.GetPath("run"), "*.far").Length; timeProvider.Advance(TimeSpan.FromSeconds(13)); await WaitUntilAsync(() => attempts == 2); await stopping.CancelAsync(); @@ -335,7 +339,9 @@ await monitor.Completion.WaitAsync( TestContext.Current.CancellationToken); // assert - Assert.Equal("True|3|2|1", $"{rolledBack}|{handler.RequestCount}|{attempts}|{adoptions}"); + Assert.Equal( + "True|1|3|2|1", + $"{rolledBack}|{filesAfterRollback}|{handler.RequestCount}|{attempts}|{adoptions}"); } private async Task CreateCoordinatorAsync( diff --git a/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/NitroExtensionsTests.cs b/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/NitroExtensionsTests.cs index d2f8941d752..818c60b5d99 100644 --- a/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/NitroExtensionsTests.cs +++ b/src/HotChocolate/Fusion/test/Fusion.Aspire.Tests/NitroExtensionsTests.cs @@ -348,7 +348,8 @@ public void AddNitro_Should_ConfigureSeedUpdates() // act builder.AddNitro( "production", - options => + portalUrl: null, + configureSeedUpdates: options => { options.Enabled = false; options.AutoUpdate = false; @@ -359,6 +360,40 @@ public void AddNitro_Should_ConfigureSeedUpdates() Assert.Equal("False|False", $"{options.Enabled}|{options.AutoUpdate}"); } + [Fact] + public void AddNitro_Should_AcceptAnExplicitNullPortalUrl() + { + // arrange + var builder = DistributedApplication.CreateBuilder(); + + // act + builder.AddNitro("production", null); + + // assert + Assert.Equal("production", GetNitroCompositionOptions(builder).Coordinator?.Stage); + } + + [Fact] + public void AddNitro_Should_UpdateAutoUpdateDefault_WhenCalledAgain() + { + // arrange + var builder = DistributedApplication.CreateBuilder(); + builder.AddNitro("production"); + + // act + builder.AddNitro( + "production", + portalUrl: null, + configureSeedUpdates: options => options.AutoUpdate = false); + + // assert + var options = GetNitroCompositionOptions(builder); + Assert.Equal( + "False|False", + $"{options.SeedUpdates.AutoUpdate}|" + + $"{options.Coordinator!.IsAutoUpdateEnabled("gateway")}"); + } + [Fact] public void NitroGateway_Should_RegisterBothAutoUpdateCommands_WhenNitroIsAddedFirst() { @@ -415,7 +450,8 @@ public void SeedUpdateService_Should_NotStartMonitor_WhenDetectionIsDisabled() var builder = DistributedApplication.CreateBuilder(); builder.AddNitro( "production", - options => options.Enabled = false); + portalUrl: null, + configureSeedUpdates: options => options.Enabled = false); var gateway = builder .AddProject("gateway", GetTestProjectFile()) .WithGraphQLSchemaComposition() @@ -443,7 +479,7 @@ public void SeedUpdateService_Should_NotStartMonitor_WhenDetectionIsDisabled() } [Fact] - public async Task AutoUpdateCommands_Should_ShowOnlyDisableCommand_WhenAutoUpdateStartsEnabled() + public async Task AutoUpdateCommands_Should_HideUntilMonitorIsReady() { // arrange var builder = DistributedApplication.CreateBuilder(); @@ -479,12 +515,19 @@ public async Task AutoUpdateCommands_Should_ShowOnlyDisableCommand_WhenAutoUpdat }; // act - var states = commands.Select(command => command.UpdateState!(context)); + var beforeStart = commands.Select(command => command.UpdateState!(context)).ToArray(); + lifetime.StopApplication(); + using var gate = new SemaphoreSlim(1, 1); + service.Start( + gateway.Resource, + "QXBpCmdhdGV3YXk", + gate, + (_, _) => Task.FromResult(true)); + var afterStart = commands.Select(command => command.UpdateState!(context)).ToArray(); // assert - Assert.Equal( - [ResourceCommandState.Enabled, ResourceCommandState.Hidden], - states); + Assert.Equal("Hidden, Hidden", string.Join(", ", beforeStart)); + Assert.Equal("Enabled, Hidden", string.Join(", ", afterStart)); } ///