diff --git a/Directory.Packages.props b/Directory.Packages.props index 96ca5f181..02d73f9e3 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -77,6 +77,7 @@ + diff --git a/examples/floci/CommunityToolkit.Aspire.Hosting.Floci.AppHost.TypeScript/apphost.mts b/examples/floci/CommunityToolkit.Aspire.Hosting.Floci.AppHost.TypeScript/apphost.mts index 4c498a58d..0a9f7ccce 100644 --- a/examples/floci/CommunityToolkit.Aspire.Hosting.Floci.AppHost.TypeScript/apphost.mts +++ b/examples/floci/CommunityToolkit.Aspire.Hosting.Floci.AppHost.TypeScript/apphost.mts @@ -1,3 +1,5 @@ +import { mkdirSync } from "node:fs"; +import { randomUUID } from "node:crypto"; import path from "node:path"; import { fileURLToPath } from "node:url"; import { createBuilder } from './.aspire/modules/aspire.mjs'; @@ -9,6 +11,13 @@ const flociAws = await builder.addFlociAws('floci-aws'); const flociAzure = await builder.addFlociAzure('floci-az'); const flociGcp = await builder.addFlociGcp('floci-gcp'); +await flociAzure.withDockerSocket(); +await flociAzure.withEnvironment('FLOCI_AZ_DOCKER_RESOURCE_NAMESPACE', `sample-${randomUUID()}`); +const serviceBus = await flociAzure.withServiceBus({ name: 'servicebus' }); +await serviceBus.amqpEndpoint(); +await serviceBus.amqpTlsEndpoint(); +await serviceBus.connectionStringExpression(); + // A single Floci UI console browses all three clouds — flociAws.withFlociUI() creates the // console wired to AWS, then withCloudReference* attaches the Azure and GCP resources to it. await flociAws.withFlociUI({ @@ -28,10 +37,42 @@ const apiService = await builder.addProject("floci-api", apiServiceProjectPath) .withFlociAwsReference(flociAws) .withFlociAzureReference(flociAzure) .withFlociGcpReference(flociGcp) + .withFlociAzureServiceBusReference(serviceBus) .waitFor(flociAws) .waitFor(flociAzure) .waitFor(flociGcp); +// Resolve the returned child handle from a real container and require an AMQP response. +await builder.addContainer('servicebus-probe', 'node:22-alpine') + .withFlociAzureServiceBusReference(serviceBus) + .withHttpEndpoint({ targetPort: 8080 }) + .withHttpHealthCheck({ path: '/' }) + .waitFor(flociAzure) + .withArgs(['-e', ` + const net = require('node:net'); + const http = require('node:http'); + const assert = require('node:assert/strict'); + const connectionString = process.env.ConnectionStrings__servicebus; + const endpoint = new URL(connectionString.split(';')[0].slice('Endpoint='.length)); + const header = Buffer.from([65, 77, 81, 80, 0, 1, 0, 0]); + let ready = false; + const socket = net.connect({ host: endpoint.hostname, port: Number(endpoint.port) }, () => socket.write(header)); + socket.setTimeout(20000, () => { throw new Error('AMQP handshake timed out'); }); + let response = Buffer.alloc(0); + socket.on('data', data => { + response = Buffer.concat([response, data]); + if (response.length >= header.length) { + assert.deepEqual(response.subarray(0, header.length), header); + ready = true; + socket.destroy(); + } + }); + http.createServer((request, response) => { + response.writeHead(ready ? 200 : 503); + response.end(); + }).listen(8080, '0.0.0.0'); + `]); + // ── Custom port and region ──────────────────────────────────────────────────── await builder.addFlociAws('floci-custom', { port: 14566, @@ -51,7 +92,9 @@ await flociPersistent.withDataVolume('floci-data'); // ── Persistent storage — bind mount ─────────────────────────────────────────── const flociMount = await builder.addFlociAws('floci-mount'); -await flociMount.withDataBindMount('/tmp/floci-data'); +const flociDataPath = path.join(appHostDirectory, '.aspire', 'floci-data'); +mkdirSync(flociDataPath, { recursive: true }); +await flociMount.withDataBindMount(flociDataPath); // ── Compile-time coverage ───────────────────────────────────────────────────── // Guards with false so these are type-checked but never executed. @@ -82,6 +125,7 @@ if (includeCompileOnlyScenarios) { await _azurePodman.withDockerSocket({ socketPath: '/run/user/1000/podman/podman.sock' }); + await _azurePodman.withServiceBus({ name: 'podman-servicebus', amqpPort: 15673, amqpTlsPort: 15674 }); const _gcpPodman = await builder.addFlociGcp('floci-gcp-podman'); await _gcpPodman.withDockerSocket({ diff --git a/src/CommunityToolkit.Aspire.Hosting.Floci/FlociAzureServiceBusConnectionString.cs b/src/CommunityToolkit.Aspire.Hosting.Floci/FlociAzureServiceBusConnectionString.cs new file mode 100644 index 000000000..860888705 --- /dev/null +++ b/src/CommunityToolkit.Aspire.Hosting.Floci/FlociAzureServiceBusConnectionString.cs @@ -0,0 +1,21 @@ +using Aspire.Hosting.ApplicationModel; + +namespace CommunityToolkit.Aspire.Hosting.Floci; + +// The sidecar publishes directly on the Docker host. Expose the parent as the dependency +// so Aspire does not create a container tunnel for the child's proxyless listeners. +internal sealed class FlociAzureServiceBusConnectionString(FlociAzureServiceBusResource resource) + : IValueProvider, IManifestExpressionProvider, IValueWithReferences +{ + internal const string ContainerHost = "floci-servicebus-host.internal"; + + public async ValueTask GetValueAsync(CancellationToken cancellationToken = default) + { + string? port = await resource.AmqpEndpoint.Property(EndpointProperty.Port).GetValueAsync(cancellationToken); + return $"Endpoint=sb://{ContainerHost}:{port};SharedAccessKeyName=RootManageSharedAccessKey;SharedAccessKey={FlociAzureServiceBusResource.DefaultSasKey};UseDevelopmentEmulator=true;"; + } + + public string ValueExpression => throw new NotSupportedException("Floci Service Bus references are only supported in run mode."); + + public IEnumerable References => [resource.Parent]; +} diff --git a/src/CommunityToolkit.Aspire.Hosting.Floci/FlociAzureServiceBusResource.cs b/src/CommunityToolkit.Aspire.Hosting.Floci/FlociAzureServiceBusResource.cs new file mode 100644 index 000000000..959bf9eef --- /dev/null +++ b/src/CommunityToolkit.Aspire.Hosting.Floci/FlociAzureServiceBusResource.cs @@ -0,0 +1,67 @@ +namespace Aspire.Hosting.ApplicationModel; + +/// +/// Represents the Service Bus AMQP data plane exposed by a Floci Azure emulator resource. +/// +/// +/// floci-az serves Service Bus AMQP from an Artemis sidecar container that publishes the +/// configured ports directly on the Docker host. The resource models those host ports as +/// proxyless Aspire endpoints so DCP can allocate them without trying to proxy traffic to the +/// parent container. Container consumers reach those same ports through the container runtime's +/// host gateway. This resource is only supported in run mode. +/// +/// The name of the resource. +/// The parent Floci Azure emulator resource. +[AspireExport(ExposeProperties = true)] +public class FlociAzureServiceBusResource( + string name, + FlociAzureContainerResource parent) : Resource(name), + IResourceWithParent, + IResourceWithConnectionString, + IResourceWithEndpoints +{ + internal const string DefaultName = "servicebus"; + internal const string AmqpEndpointName = "amqp"; + internal const string AmqpTlsEndpointName = "amqps"; + + private EndpointReference? _amqpEndpoint; + private EndpointReference? _amqpTlsEndpoint; + + // Placeholder from the official Service Bus emulator's connection-string shape; floci-az + // does not enforce authentication, the SDK only requires the component to be present. + internal const string DefaultSasKey = "SAS_KEY_VALUE"; + + /// + /// Gets the parent Floci Azure emulator resource. + /// + public FlociAzureContainerResource Parent { get; } = parent ?? throw new ArgumentNullException(nameof(parent)); + + /// + /// Gets the Service Bus plain AMQP endpoint. + /// + public EndpointReference AmqpEndpoint => + _amqpEndpoint ??= new EndpointReference(this, AmqpEndpointName); + + /// + /// Gets the Service Bus AMQPS/TLS endpoint. + /// + public EndpointReference AmqpTlsEndpoint => + _amqpTlsEndpoint ??= new EndpointReference(this, AmqpTlsEndpointName); + + /// + /// Gets the Service Bus connection string expression. + /// UseDevelopmentEmulator=true makes the Azure SDKs use plain AMQP (no TLS), matching + /// the official Service Bus emulator's connection-string shape. + /// + public ReferenceExpression ConnectionStringExpression => + ReferenceExpression.Create( + $"Endpoint={AmqpEndpoint};SharedAccessKeyName=RootManageSharedAccessKey;SharedAccessKey={DefaultSasKey};UseDevelopmentEmulator=true;"); + + IEnumerable> IResourceWithConnectionString.GetConnectionProperties() => + Parent.CombineProperties([ + new("Host", ReferenceExpression.Create($"{AmqpEndpoint.Property(EndpointProperty.Host)}")), + new("Port", ReferenceExpression.Create($"{AmqpEndpoint.Property(EndpointProperty.Port)}")), + new("Uri", ReferenceExpression.Create($"{AmqpEndpoint}")), + new("Endpoint", ReferenceExpression.Create($"{AmqpEndpoint}")) + ]); +} diff --git a/src/CommunityToolkit.Aspire.Hosting.Floci/FlociHostingExtension.Azure.cs b/src/CommunityToolkit.Aspire.Hosting.Floci/FlociHostingExtension.Azure.cs index 90ec94c08..a633c45e8 100644 --- a/src/CommunityToolkit.Aspire.Hosting.Floci/FlociHostingExtension.Azure.cs +++ b/src/CommunityToolkit.Aspire.Hosting.Floci/FlociHostingExtension.Azure.cs @@ -1,3 +1,4 @@ +using System.Globalization; using Aspire.Hosting.ApplicationModel; using CommunityToolkit.Aspire.Hosting.Floci; @@ -73,6 +74,146 @@ public static IResourceBuilder WithReference( }); + /// + /// Adds a child resource representing the Service Bus AMQP data plane exposed by the Floci + /// Azure emulator, enabling the data plane on the emulator (MOCKED=false, + /// START_ON_BOOT=true). + /// + /// + /// Reference the returned resource with Aspire's standard WithReference API to inject + /// its Service Bus connection string (e.g. for AddAzureServiceBusClient). The Artemis + /// sidecar is a separate container floci-az starts via Docker, so the emulator also needs + /// . + /// Aspire allocates proxyless AMQP host endpoints when ports are not specified. Requires + /// floci-az 0.12.0 or later. This API is run-only; call it inside an + /// ExecutionContext.IsRunMode branch when publishing the AppHost. + /// + /// Adds a Service Bus child resource to the Floci Azure emulator + /// The Floci Azure resource builder. + /// The globally unique name of the Service Bus resource (default: servicebus). + /// Pass a distinct name for each additional emulator. + /// Host port for plain AMQP (default: allocated by Aspire). + /// Host port for AMQPS/TLS (default: allocated by Aspire). + /// A reference to the for further configuration. + [AspireExport] + public static IResourceBuilder WithServiceBus( + this IResourceBuilder builder, + [ResourceName] string name = FlociAzureServiceBusResource.DefaultName, + int? amqpPort = null, + int? amqpTlsPort = null) + { + ArgumentNullException.ThrowIfNull(builder); + ArgumentException.ThrowIfNullOrWhiteSpace(name); + + if (builder.ApplicationBuilder.ExecutionContext.IsPublishMode) + { + throw new NotSupportedException( + "Floci Service Bus is only supported in run mode. Call WithServiceBus inside an ExecutionContext.IsRunMode branch and reference a deployable Service Bus resource when publishing."); + } + + if (amqpPort is not null && amqpPort == amqpTlsPort) + { + throw new ArgumentException("AMQP and AMQPS must use different host ports.", nameof(amqpTlsPort)); + } + + FlociAzureServiceBusResource? existing = builder.ApplicationBuilder.Resources + .OfType() + .FirstOrDefault(resource => resource.Parent == builder.Resource); + if (existing is not null) + { + if (!string.Equals(name, existing.Name, StringComparison.OrdinalIgnoreCase)) + { + throw new InvalidOperationException( + $"Service Bus is already configured on '{builder.Resource.Name}' as '{existing.Name}' and cannot be reconfigured with a different name."); + } + + int? existingAmqpPort = existing.AmqpEndpoint.EndpointAnnotation.Port; + int? existingAmqpTlsPort = existing.AmqpTlsEndpoint.EndpointAnnotation.Port; + if ((amqpPort is not null && amqpPort != existingAmqpPort) + || (amqpTlsPort is not null && amqpTlsPort != existingAmqpTlsPort)) + { + throw new InvalidOperationException( + $"Service Bus is already configured on '{builder.Resource.Name}' and cannot be reconfigured with different ports."); + } + + return builder.ApplicationBuilder.CreateResourceBuilder(existing); + } + + var serviceBus = new FlociAzureServiceBusResource(name, builder.Resource); + var serviceBusBuilder = builder.ApplicationBuilder + .AddResource(serviceBus) + .WithEndpoint( + port: amqpPort, + scheme: "sb", + name: FlociAzureServiceBusResource.AmqpEndpointName, + isProxied: false) + .WithEndpoint( + port: amqpTlsPort, + scheme: "amqps", + name: FlociAzureServiceBusResource.AmqpTlsEndpointName, + isProxied: false) + .WithParentRelationship(builder) + .ExcludeFromManifest(); + + builder.WithEnvironment(context => + { + if (context.ExecutionContext.IsPublishMode) + { + return; + } + + context.EnvironmentVariables["FLOCI_AZ_SERVICES_SERVICE_BUS_MOCKED"] = "false"; + context.EnvironmentVariables["FLOCI_AZ_SERVICES_SERVICE_BUS_START_ON_BOOT"] = "true"; + context.EnvironmentVariables["FLOCI_AZ_SERVICES_SERVICE_BUS_AMQP_PORT"] = + serviceBus.AmqpEndpoint.Port.ToString(CultureInfo.InvariantCulture); + context.EnvironmentVariables["FLOCI_AZ_SERVICES_SERVICE_BUS_AMQP_TLS_PORT"] = + serviceBus.AmqpTlsEndpoint.Port.ToString(CultureInfo.InvariantCulture); + }); + + return serviceBusBuilder; + } + + /// + /// Adds a Service Bus connection string reference, using the Docker host gateway for container consumers. + /// + /// Adds a Floci Azure Service Bus reference + /// The type of the resource receiving the reference. + /// The resource builder receiving the reference. + /// The Service Bus child resource to reference. + /// The connection string name, defaulting to the child resource name. + /// The resource builder receiving the reference. + [AspireExport("withFlociAzureServiceBusReference")] + public static IResourceBuilder WithReference( + this IResourceBuilder builder, + IResourceBuilder serviceBus, + string? connectionName = null) + where TDestination : IResourceWithEnvironment + { + ArgumentNullException.ThrowIfNull(builder); + ArgumentNullException.ThrowIfNull(serviceBus); + + if (builder.ApplicationBuilder.ExecutionContext.IsPublishMode) + { + throw new NotSupportedException("Floci Service Bus references are only supported in run mode."); + } + + if (builder.Resource is not ContainerResource container) + { + return ResourceBuilderExtensions.WithReference(builder, serviceBus, connectionName); + } + + builder.ApplicationBuilder.CreateResourceBuilder(container) + .WithContainerRuntimeArgs("--add-host", $"{FlociAzureServiceBusConnectionString.ContainerHost}:host-gateway"); + + return builder + .WithEnvironment(context => + { + context.EnvironmentVariables[$"ConnectionStrings__{connectionName ?? serviceBus.Resource.Name}"] = + new FlociAzureServiceBusConnectionString(serviceBus.Resource); + }) + .WithRelationship(serviceBus.Resource, "Reference"); + } + /// /// Adds a child resource representing the Cosmos DB API exposed by the Floci Azure emulator. /// diff --git a/src/CommunityToolkit.Aspire.Hosting.Floci/README.md b/src/CommunityToolkit.Aspire.Hosting.Floci/README.md index 7c35d914e..4d5b53d42 100644 --- a/src/CommunityToolkit.Aspire.Hosting.Floci/README.md +++ b/src/CommunityToolkit.Aspire.Hosting.Floci/README.md @@ -107,6 +107,46 @@ builder.AddAzureCosmosClient("cosmos"); The Cosmos child resource is additive, so combine `WithReference(cosmos)` with `WithReference(azure)` when you also want the base endpoint / storage variables. (Talking to the floci Cosmos emulator over HTTP from the .NET SDK still needs the usual client-side settings — Gateway mode, and HTTP/1.1 — which are the app's concern, as with any local Cosmos emulator.) +For **Service Bus**, use `WithServiceBus()` / `withServiceBus()` to model the AMQP data plane as a child resource, then reference it with `WithReference()` / `withFlociAzureServiceBusReference()`: + +```csharp +var azure = builder.AddFlociAzure("floci-az") + .WithDockerSocket(); +var serviceBus = azure.WithServiceBus(); + +builder.AddProject("api") + .WithReference(serviceBus) // ConnectionStrings__servicebus + .WaitFor(azure); +``` + +```typescript +const azure = (await builder.addFlociAzure('floci-az')).withDockerSocket(); +const serviceBus = await azure.withServiceBus(); + +await builder.addProject('api', '../MyApi/MyApi.csproj') + .withFlociAzureServiceBusReference(serviceBus) + .waitFor(azure); +``` + +App side, this is the standard Aspire flow: + +```csharp +builder.AddAzureServiceBusClient("servicebus"); +``` + +| Variable | Value | +|---|---| +| `ConnectionStrings__{resourceName}` (default `servicebus`) | `Endpoint=sb://{host}:{amqpPort};SharedAccessKeyName=RootManageSharedAccessKey;SharedAccessKey=SAS_KEY_VALUE;UseDevelopmentEmulator=true;` — `{host}` is `localhost` for host processes and the container runtime's host gateway for containers | + +`WithServiceBus` sets `FLOCI_AZ_SERVICES_SERVICE_BUS_MOCKED=false` and `FLOCI_AZ_SERVICES_SERVICE_BUS_START_ON_BOOT=true`. Aspire models the sidecar's AMQP and AMQPS host ports as proxyless endpoints and allocates them by default; pass `amqpPort` / `amqpTlsPort` to use fixed ports. The management plane (for example, `ServiceBusAdministrationClient`) remains on the base endpoint from `WithReference(azure)`. Requires `WithDockerSocket()` and floci-az 0.12.0 or later. + +AMQP and AMQPS must use different host ports. Repeated calls on the same emulator return the existing child only when the resource name and any specified ports match. + +For container consumers, the Service Bus reference helper adds `floci-servicebus-host.internal:host-gateway` to the container's host mappings and uses the sidecar's published AMQP port. This requires a container runtime supporting Docker's `--add-host=...:host-gateway` option. Use the Service Bus reference helper instead of passing the child's raw endpoint or connection string expression to a container. + +Service Bus support is run-only. `WithServiceBus()` throws in publish mode because the Artemis sidecar has no deployable backing resource in the Aspire model. In an AppHost that also publishes, call it inside `if (builder.ExecutionContext.IsRunMode)` and reference a deployable Service Bus resource in the publish branch. + +When running multiple Floci Azure emulators on the same Docker host, set a different `FLOCI_AZ_DOCKER_RESOURCE_NAMESPACE` environment variable on each emulator to keep their sidecar container names separate. Also pass a globally unique child name to each additional `WithServiceBus` call, for example `secondAzure.WithServiceBus("second-servicebus")`. The first can keep the default `servicebus` connection name; references to the second inject `ConnectionStrings__second-servicebus`. Port allocation alone does not isolate resource or sidecar names. **GCP** diff --git a/tests/CommunityToolkit.Aspire.Hosting.Floci.Tests/AzureServiceBusFunctionalTests.cs b/tests/CommunityToolkit.Aspire.Hosting.Floci.Tests/AzureServiceBusFunctionalTests.cs new file mode 100644 index 000000000..8bdea7367 --- /dev/null +++ b/tests/CommunityToolkit.Aspire.Hosting.Floci.Tests/AzureServiceBusFunctionalTests.cs @@ -0,0 +1,116 @@ +using Aspire.Hosting; +using Aspire.Hosting.Utils; +using Azure.Messaging.ServiceBus; +using Azure.Messaging.ServiceBus.Administration; +using CommunityToolkit.Aspire.Testing; +using System.Net.Sockets; + +namespace CommunityToolkit.Aspire.Hosting.Floci.Tests; + +[RequiresDocker] +public class AzureServiceBusFunctionalTests(ITestOutputHelper testOutputHelper) +{ + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task ServiceBusStartsAndSupportsSdkSendAndReceive(bool useBatch) + { + using var timeout = CancellationTokenSource.CreateLinkedTokenSource(TestContext.Current.CancellationToken); + timeout.CancelAfter(TimeSpan.FromMinutes(3)); + var cancellationToken = timeout.Token; + + using var builder = TestDistributedApplicationBuilder.Create(testOutputHelper); + string resourceNamespace = $"test-{Guid.NewGuid():N}"; + var azure = builder.AddFlociAzure("floci-az") + .WithDockerSocket() + .WithEnvironment("FLOCI_AZ_DOCKER_RESOURCE_NAMESPACE", resourceNamespace); + var serviceBus = azure.WithServiceBus(); + var consumer = builder.AddExecutable("consumer", "dotnet", builder.AppHostDirectory, "--version") + .WithReference(serviceBus) + .WaitFor(azure); + var containerConsumer = builder.AddContainer("container-consumer", "node", "22-alpine") + .WithReference(serviceBus) + .WaitFor(azure) + .WithArgs("-e", """ + const net = require('node:net'); + const assert = require('node:assert/strict'); + const connectionString = process.env.ConnectionStrings__servicebus; + const endpoint = new URL(connectionString.split(';')[0].slice('Endpoint='.length)); + assert.notEqual(endpoint.hostname, 'localhost'); + const header = Buffer.from([65, 77, 81, 80, 0, 1, 0, 0]); + const socket = net.connect({ host: endpoint.hostname, port: Number(endpoint.port) }, () => socket.write(header)); + socket.setTimeout(20000, () => { throw new Error('AMQP handshake timed out'); }); + let response = Buffer.alloc(0); + socket.on('data', data => { + response = Buffer.concat([response, data]); + if (response.length >= header.length) { + assert.deepEqual(response.subarray(0, header.length), header); + socket.destroy(); + } + }); + socket.on('end', () => assert.ok(response.length >= header.length, 'Missing AMQP response')); + """); + + await using var app = await builder.BuildAsync(cancellationToken); + await app.StartAsync(cancellationToken); + await app.ResourceNotifications.WaitForResourceHealthyAsync(azure.Resource.Name, cancellationToken); + + await foreach (var resourceEvent in app.ResourceNotifications.WatchAsync(cancellationToken)) + { + if (resourceEvent.Resource == containerConsumer.Resource && resourceEvent.Snapshot.ExitCode is { } exitCode) + { + Assert.Equal(0, exitCode); + break; + } + } + + Assert.NotEqual(serviceBus.Resource.AmqpEndpoint.Port, serviceBus.Resource.AmqpTlsEndpoint.Port); + + // Both listeners must exist before any management call can lazily start the namespace. + foreach (var endpoint in new[] { serviceBus.Resource.AmqpEndpoint, serviceBus.Resource.AmqpTlsEndpoint }) + { + using var socket = new TcpClient(); + await socket.ConnectAsync(endpoint.Host, endpoint.Port, cancellationToken); + } + + var environment = await consumer.Resource.GetEnvironmentVariablesAsync(serviceProvider: app.Services) + .AsTask().WaitAsync(cancellationToken); + string connectionString = environment["ConnectionStrings__servicebus"]; + Assert.Equal(await app.GetConnectionStringAsync(serviceBus.Resource.Name, cancellationToken), connectionString); + var managementEndpoint = new Uri(azure.Resource.PrimaryEndpoint.Url); + var administration = new ServiceBusAdministrationClient( + $"Endpoint=sb://{managementEndpoint.Authority};SharedAccessKeyName=RootManageSharedAccessKey;" + + "SharedAccessKey=SAS_KEY_VALUE;UseDevelopmentEmulator=true;"); + string queueName = $"test-{Guid.NewGuid():N}"; + await administration.CreateQueueAsync(queueName, cancellationToken); + + await using var client = new ServiceBusClient(connectionString); + await using var sender = client.CreateSender(queueName); + string[] bodies = useBatch ? ["first", "second", "third"] : ["single"]; + if (useBatch) + { + using var batch = await sender.CreateMessageBatchAsync(cancellationToken); + foreach (string body in bodies) + { + Assert.True(batch.TryAddMessage(new ServiceBusMessage(body))); + } + + await sender.SendMessagesAsync(batch, cancellationToken); + } + else + { + await sender.SendMessageAsync(new ServiceBusMessage(bodies[0]), cancellationToken); + } + + await using var receiver = client.CreateReceiver(queueName); + foreach (string body in bodies) + { + var message = await receiver.ReceiveMessageAsync(TimeSpan.FromSeconds(20), cancellationToken); + Assert.NotNull(message); + Assert.Equal(body, message.Body.ToString()); + await receiver.CompleteMessageAsync(message, cancellationToken); + } + + await administration.DeleteQueueAsync(queueName, cancellationToken); + } +} diff --git a/tests/CommunityToolkit.Aspire.Hosting.Floci.Tests/AzureServiceBusResourceTests.cs b/tests/CommunityToolkit.Aspire.Hosting.Floci.Tests/AzureServiceBusResourceTests.cs new file mode 100644 index 000000000..5dbc2de26 --- /dev/null +++ b/tests/CommunityToolkit.Aspire.Hosting.Floci.Tests/AzureServiceBusResourceTests.cs @@ -0,0 +1,280 @@ +using Aspire.Hosting; +using Aspire.Hosting.Utils; +using CommunityToolkit.Aspire.Testing; + +namespace CommunityToolkit.Aspire.Hosting.Floci.Tests; + +public class AzureServiceBusResourceTests +{ + [Fact] + public void WithServiceBusCreatesChildResourceWithAspireAllocatedEndpoints() + { + IDistributedApplicationBuilder builder = DistributedApplication.CreateBuilder(); + + var azure = builder.AddFlociAzure("floci-az"); + azure.WithServiceBus(); + + using var app = builder.Build(); + var appModel = app.Services.GetRequiredService(); + + var serviceBus = Assert.Single(appModel.Resources.OfType()); + Assert.Equal("servicebus", serviceBus.Name); + Assert.Same(azure.Resource, serviceBus.Parent); + + Assert.Collection( + serviceBus.Annotations.OfType().OrderBy(endpoint => endpoint.Name), + amqp => AssertEndpoint(amqp, "amqp", "sb", null), + amqps => AssertEndpoint(amqps, "amqps", "amqps", null)); + } + + [Fact] + public void WithServiceBusHonorsExplicitPorts() + { + IDistributedApplicationBuilder builder = DistributedApplication.CreateBuilder(); + + var serviceBus = builder.AddFlociAzure("floci-az") + .WithServiceBus(amqpPort: 5673, amqpTlsPort: 5674); + + Assert.Equal(5673, serviceBus.Resource.AmqpEndpoint.EndpointAnnotation.Port); + Assert.Equal(5674, serviceBus.Resource.AmqpTlsEndpoint.EndpointAnnotation.Port); + } + + [Fact] + public void EqualListenerPortsThrowBeforeRegisteringChild() + { + IDistributedApplicationBuilder builder = DistributedApplication.CreateBuilder(); + var azure = builder.AddFlociAzure("floci-az"); + + var exception = Assert.Throws(() => azure.WithServiceBus(amqpPort: 5673, amqpTlsPort: 5673)); + + Assert.Equal("amqpTlsPort", exception.ParamName); + Assert.DoesNotContain(builder.Resources, resource => resource is FlociAzureServiceBusResource); + } + + [Fact] + public void PublishModeRejectsServiceBusBeforeRegisteringChild() + { + using var builder = TestDistributedApplicationBuilder.Create(DistributedApplicationOperation.Publish); + var azure = builder.AddFlociAzure("floci-az"); + + var exception = Assert.Throws(() => azure.WithServiceBus()); + + Assert.Contains("ExecutionContext.IsRunMode", exception.Message); + Assert.DoesNotContain(builder.Resources, resource => resource is FlociAzureServiceBusResource); + } + + [Fact] + public async Task WithServiceBusUsesAllocatedEndpointPortsForTheDataPlane() + { + IDistributedApplicationBuilder builder = DistributedApplication.CreateBuilder(); + + var azure = builder.AddFlociAzure("floci-az"); + var serviceBus = azure.WithServiceBus(); + AllocateEndpoints(serviceBus.Resource, 5673, 5674); + + using var app = builder.Build(); + var appModel = app.Services.GetRequiredService(); + + var resource = Assert.Single(appModel.Resources.OfType()); + Assert.True(resource.TryGetAnnotationsOfType(out IEnumerable? envAnnotations)); + + var envVars = new Dictionary(); + var context = new EnvironmentCallbackContext(builder.ExecutionContext, envVars); + foreach (var annotation in envAnnotations!) + { + await annotation.Callback(context); + } + + Assert.Equal("false", envVars["FLOCI_AZ_SERVICES_SERVICE_BUS_MOCKED"]); + Assert.Equal("true", envVars["FLOCI_AZ_SERVICES_SERVICE_BUS_START_ON_BOOT"]); + Assert.Equal("5673", envVars["FLOCI_AZ_SERVICES_SERVICE_BUS_AMQP_PORT"]); + Assert.Equal("5674", envVars["FLOCI_AZ_SERVICES_SERVICE_BUS_AMQP_TLS_PORT"]); + } + + [Fact] + public async Task WithServiceBusDoesNotConfigureTheDataPlaneInPublishMode() + { + IDistributedApplicationBuilder builder = DistributedApplication.CreateBuilder(); + + var azure = builder.AddFlociAzure("floci-az"); + azure.WithServiceBus(); + + var envVars = new Dictionary(); + var executionContext = new DistributedApplicationExecutionContext( + new DistributedApplicationExecutionContextOptions(DistributedApplicationOperation.Publish)); + var context = new EnvironmentCallbackContext(executionContext, envVars); + + foreach (var annotation in azure.Resource.Annotations.OfType()) + { + await annotation.Callback(context); + } + + Assert.DoesNotContain("FLOCI_AZ_SERVICES_SERVICE_BUS_MOCKED", envVars); + Assert.DoesNotContain("FLOCI_AZ_SERVICES_SERVICE_BUS_START_ON_BOOT", envVars); + Assert.DoesNotContain("FLOCI_AZ_SERVICES_SERVICE_BUS_AMQP_PORT", envVars); + Assert.DoesNotContain("FLOCI_AZ_SERVICES_SERVICE_BUS_AMQP_TLS_PORT", envVars); + } + + [Fact] + public async Task ConnectionStringMatchesTheEmulatorShape() + { + IDistributedApplicationBuilder builder = DistributedApplication.CreateBuilder(); + + var serviceBus = builder.AddFlociAzure("floci-az").WithServiceBus(); + AllocateEndpoints(serviceBus.Resource, 5673, 5674); + + string? connectionString = await serviceBus.Resource.ConnectionStringExpression + .GetValueAsync(CancellationToken.None); + + Assert.Equal( + "Endpoint=sb://localhost:5673;SharedAccessKeyName=RootManageSharedAccessKey;" + + "SharedAccessKey=SAS_KEY_VALUE;UseDevelopmentEmulator=true;", + connectionString); + } + + [Fact] + public void SecondWithServiceBusReturnsTheExistingChild() + { + IDistributedApplicationBuilder builder = DistributedApplication.CreateBuilder(); + + var azure = builder.AddFlociAzure("floci-az"); + var first = azure.WithServiceBus(amqpPort: 5673); + var second = azure.WithServiceBus(amqpPort: 5673); + + Assert.Same(first.Resource, second.Resource); + + using var app = builder.Build(); + var appModel = app.Services.GetRequiredService(); + Assert.Single(appModel.Resources.OfType()); + } + + [Fact] + public void ConflictingPortsThrow() + { + IDistributedApplicationBuilder builder = DistributedApplication.CreateBuilder(); + + var azure = builder.AddFlociAzure("floci-az"); + azure.WithServiceBus(amqpPort: 5673); + + Assert.Throws(() => azure.WithServiceBus(amqpPort: 5675)); + } + + [Fact] + public void ConflictingNamesThrow() + { + IDistributedApplicationBuilder builder = DistributedApplication.CreateBuilder(); + var azure = builder.AddFlociAzure("floci-az"); + var serviceBus = azure.WithServiceBus("first"); + + Assert.Throws(() => azure.WithServiceBus("second")); + Assert.Same(serviceBus.Resource, azure.WithServiceBus("FIRST").Resource); + Assert.Single(builder.Resources.OfType()); + } + + [Fact] + public async Task MultipleEmulatorsUseDistinctChildNamesAndConnectionKeys() + { + IDistributedApplicationBuilder builder = DistributedApplication.CreateBuilder(); + var first = builder.AddFlociAzure("first").WithServiceBus(); + var second = builder.AddFlociAzure("second").WithServiceBus("second-servicebus"); + AllocateEndpoints(first.Resource, 5673, 5674); + AllocateEndpoints(second.Resource, 5675, 5676); + var consumer = builder.AddExecutable("consumer", "dotnet", ".") + .WithReference(first) + .WithReference(second); + + using var app = builder.Build(); + var environment = await consumer.Resource.GetEnvironmentVariablesAsync(serviceProvider: app.Services); + + Assert.Contains("Endpoint=sb://localhost:5673;", environment["ConnectionStrings__servicebus"]); + Assert.Contains("Endpoint=sb://localhost:5675;", environment["ConnectionStrings__second-servicebus"]); + Assert.Equal(2, builder.Resources.OfType().Count()); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task WithReferenceInjectsTheConnectionString(bool useContainer) + { + IDistributedApplicationBuilder builder = DistributedApplication.CreateBuilder(); + + var serviceBus = builder.AddFlociAzure("floci-az").WithServiceBus(); + AllocateEndpoints(serviceBus.Resource, 5673, 5674); + + IResourceBuilder consumer = useContainer + ? builder.AddContainer("api", "my-api-image").WithReference(serviceBus) + : builder.AddExecutable("api", "dotnet", ".").WithReference(serviceBus); + + using var app = builder.Build(); + + Assert.True(consumer.Resource.TryGetAnnotationsOfType( + out IEnumerable? envAnnotations)); + + var envVars = new Dictionary(); + var context = new EnvironmentCallbackContext(builder.ExecutionContext, envVars); + foreach (var annotation in envAnnotations!) + { + await annotation.Callback(context); + } + + object connectionString = envVars["ConnectionStrings__servicebus"]; + string? value = connectionString is IValueProvider provider + ? await provider.GetValueAsync(CancellationToken.None) + : connectionString.ToString(); + + Assert.Equal( + $"Endpoint=sb://{(useContainer ? "floci-servicebus-host.internal" : "localhost")}:5673;SharedAccessKeyName=RootManageSharedAccessKey;" + + "SharedAccessKey=SAS_KEY_VALUE;UseDevelopmentEmulator=true;", + value); + } + + [Fact] + public async Task ContainerReferencePreservesCustomConnectionNameAndAvoidsHostTunnelDependency() + { + using var builder = TestDistributedApplicationBuilder.Create(); + var azure = builder.AddFlociAzure("floci-az"); + var serviceBus = azure.WithServiceBus(); + var consumer = builder.AddContainer("api", "my-api-image") + .WithReference(serviceBus, connectionName: "messages"); + + // Dependency discovery runs before allocation and must not resolve the host ports + // or request an Aspire tunnel for the sidecar's host-only endpoints. + var dependencies = await consumer.Resource.GetResourceDependenciesAsync(builder.ExecutionContext, + new ResourceDependencyDiscoveryOptions { DiscoveryMode = ResourceDependencyDiscoveryMode.DirectOnly }); + Assert.Contains(azure.Resource, dependencies); + Assert.DoesNotContain(serviceBus.Resource, dependencies); + + AllocateEndpoints(serviceBus.Resource, 5673, 5674); + await using var app = await builder.BuildAsync(); + var environment = await consumer.Resource.GetEnvironmentVariablesAsync(serviceProvider: app.Services); + Assert.Contains("Endpoint=sb://floci-servicebus-host.internal:5673;", environment["ConnectionStrings__messages"]); + Assert.DoesNotContain("ConnectionStrings__servicebus", environment); + + var exception = await Assert.ThrowsAsync(async () => + await consumer.Resource.GetEnvironmentVariablesAsync(DistributedApplicationOperation.Publish, app.Services)); + Assert.IsType(Assert.Single(exception.InnerExceptions)); + } + + private static void AssertEndpoint( + EndpointAnnotation endpoint, + string name, + string scheme, + int? port) + { + Assert.Equal(name, endpoint.Name); + Assert.Equal(scheme, endpoint.UriScheme); + Assert.Equal(port, endpoint.Port); + Assert.False(endpoint.IsExplicitlyProxied); + } + + private static void AllocateEndpoints( + FlociAzureServiceBusResource serviceBus, + int amqpPort, + int amqpTlsPort) + { + serviceBus.AmqpEndpoint.EndpointAnnotation.AllocatedEndpoint = + new AllocatedEndpoint(serviceBus.AmqpEndpoint.EndpointAnnotation, "localhost", amqpPort); + serviceBus.AmqpTlsEndpoint.EndpointAnnotation.AllocatedEndpoint = + new AllocatedEndpoint(serviceBus.AmqpTlsEndpoint.EndpointAnnotation, "localhost", amqpTlsPort); + } +} diff --git a/tests/CommunityToolkit.Aspire.Hosting.Floci.Tests/CommunityToolkit.Aspire.Hosting.Floci.Tests.csproj b/tests/CommunityToolkit.Aspire.Hosting.Floci.Tests/CommunityToolkit.Aspire.Hosting.Floci.Tests.csproj index f408eb6d6..da8f129e8 100644 --- a/tests/CommunityToolkit.Aspire.Hosting.Floci.Tests/CommunityToolkit.Aspire.Hosting.Floci.Tests.csproj +++ b/tests/CommunityToolkit.Aspire.Hosting.Floci.Tests/CommunityToolkit.Aspire.Hosting.Floci.Tests.csproj @@ -13,6 +13,7 @@ + diff --git a/tests/CommunityToolkit.Aspire.Hosting.Floci.Tests/TypeScriptAppHostTests.cs b/tests/CommunityToolkit.Aspire.Hosting.Floci.Tests/TypeScriptAppHostTests.cs index b63854c2d..f19eb4b10 100644 --- a/tests/CommunityToolkit.Aspire.Hosting.Floci.Tests/TypeScriptAppHostTests.cs +++ b/tests/CommunityToolkit.Aspire.Hosting.Floci.Tests/TypeScriptAppHostTests.cs @@ -13,7 +13,7 @@ await TypeScriptAppHostTest.Run( appHostProject: "CommunityToolkit.Aspire.Hosting.Floci.AppHost.TypeScript", packageName: "CommunityToolkit.Aspire.Hosting.Floci", exampleName: "floci", - waitForResources: ["floci-aws", "floci-az", "floci-gcp", "floci-aws-ui", "floci-custom", "floci-gcp-custom", "floci-persistent", "floci-mount"], + waitForResources: ["floci-aws", "floci-az", "floci-gcp", "floci-aws-ui", "floci-custom", "floci-gcp-custom", "floci-persistent", "floci-mount", "servicebus-probe"], cancellationToken: TestContext.Current.CancellationToken); } }