using MassTransit; using MassTransit.EntityFrameworkCoreIntegration; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using StackExchange.Redis; using System.Diagnostics; using System.Net; using System.Net.Http.Headers; using System.Net.Sockets; using System.Text.Json; using Tiku.Application; using Tiku.Application.Jobs; using Tiku.Contracts; using Tiku.Domain.Operations; using Tiku.Domain.Tenancy; using Tiku.Infrastructure; using Tiku.Infrastructure.Messaging; using Tiku.Infrastructure.Persistence; namespace Tiku.IntegrationTests; public sealed class MassTransitOutboxTests { [Fact] public async Task Bus_outbox_drains_after_real_broker_restart() { var rabbitMqHost = Environment.GetEnvironmentVariable("TIKU_TEST_RABBITMQ"); var containerName = Environment.GetEnvironmentVariable("TIKU_TEST_RABBITMQ_CONTAINER"); if (string.IsNullOrWhiteSpace(rabbitMqHost) || Environment.GetEnvironmentVariable("TIKU_TEST_RABBITMQ_RESTART") != "1" || string.IsNullOrWhiteSpace(containerName)) { return; } await using var factory = CreateRabbitFactory(rabbitMqHost); using var client = factory.CreateClient(); Assert.True(await WaitForReadyAsync(client), "API dependencies did not become ready before restart drill."); await RunDockerAsync("stop", containerName); try { Assert.True(await WaitForBrokerPortClosedAsync(new Uri(rabbitMqHost)), "RabbitMQ AMQP port remained reachable after stopping the test container."); using (var scope = factory.CreateSystemScope("Commit outbox while RabbitMQ is stopped")) { var dbContext = scope.ServiceProvider.GetRequiredService(); var publisher = scope.ServiceProvider.GetRequiredService(); await using var transaction = await dbContext.Database.BeginTransactionAsync(); await publisher.AuthorizationChangedAsync( null, null, "broker_restart_test", 3, Guid.NewGuid().ToString("N")); await dbContext.SaveChangesAsync(); await transaction.CommitAsync(); } using var verification = factory.CreateSystemScope("Verify restart outbox backlog"); var verificationDbContext = verification.ServiceProvider.GetRequiredService(); Assert.NotEmpty(await verificationDbContext.Set().ToArrayAsync()); } finally { await RunDockerAsync("start", containerName); } Assert.True(await WaitForReadyAsync(client, 240), "RabbitMQ did not become ready within 60 seconds after restart."); var drained = false; for (var attempt = 0; attempt < 240; attempt++) { using var verification = factory.CreateSystemScope("Wait for post-restart outbox drain"); var dbContext = verification.ServiceProvider.GetRequiredService(); if (!await dbContext.Set().AnyAsync()) { drained = true; break; } await Task.Delay(250); } Assert.True(drained, "Bus outbox did not drain within 60 seconds after RabbitMQ restart."); } private static async Task RunDockerAsync(string operation, string containerName) { using var process = Process.Start(new ProcessStartInfo { FileName = "docker", ArgumentList = { operation, containerName }, RedirectStandardOutput = true, RedirectStandardError = true, UseShellExecute = false }) ?? throw new InvalidOperationException("Failed to start Docker CLI for RabbitMQ restart drill."); var standardOutput = await process.StandardOutput.ReadToEndAsync(); var standardError = await process.StandardError.ReadToEndAsync(); await process.WaitForExitAsync(); if (process.ExitCode != 0) { throw new InvalidOperationException( $"docker {operation} failed for the RabbitMQ test container: {standardError}{standardOutput}"); } } private static async Task WaitForBrokerPortClosedAsync(Uri broker) { var port = broker.IsDefaultPort ? 5672 : broker.Port; for (var attempt = 0; attempt < 40; attempt++) { using var tcpClient = new TcpClient(); try { await tcpClient.ConnectAsync(broker.Host, port).WaitAsync(TimeSpan.FromMilliseconds(250)); } catch (Exception exception) when (exception is SocketException or TimeoutException) { return true; } await Task.Delay(250); } return false; } [Fact] public async Task Background_job_request_is_transactional_consumed_once_and_keeps_database_status_view() { var rabbitMqHost = Environment.GetEnvironmentVariable("TIKU_TEST_RABBITMQ"); if (string.IsNullOrWhiteSpace(rabbitMqHost)) { return; } await using var factory = CreateRabbitFactory(rabbitMqHost); using var client = factory.CreateClient(); Assert.True(await WaitForReadyAsync(client), "API dependencies did not become ready within 10 seconds."); var tenantId = Guid.NewGuid(); await factory.SeedAsync(new Tenant { Id = tenantId, Slug = tenantId.ToString("N"), Name = "RabbitMQ Background Job Tenant" }); var workerBuilder = Host.CreateApplicationBuilder(); workerBuilder.Services.AddApplication(); workerBuilder.Services.AddInfrastructure(factory.DatabaseConnectionString); workerBuilder.Services.AddReliableMessaging(CreateRabbitOptions(rabbitMqHost, configureConsumers: true)); using var worker = workerBuilder.Build(); await worker.StartAsync(); Guid rolledBackJobId; using (var scope = factory.CreateSystemScope("Roll back background job request")) { var dbContext = scope.ServiceProvider.GetRequiredService(); var jobs = scope.ServiceProvider.GetRequiredService(); await using var transaction = await dbContext.Database.BeginTransactionAsync(); var job = await jobs.EnqueueAsync(new CreateBackgroundJobCommand( tenantId, "tenant_domain_recheck", JsonSerializer.SerializeToElement(new { }))); rolledBackJobId = job.Id; await transaction.RollbackAsync(); } using (var verification = factory.CreateSystemScope("Verify rolled back background job request")) { var dbContext = verification.ServiceProvider.GetRequiredService(); Assert.False(await dbContext.BackgroundJobs.AnyAsync(item => item.Id == rolledBackJobId)); Assert.False(await dbContext.Set().AnyAsync( item => item.MessageId == rolledBackJobId)); } BackgroundJobItem committedJob; using (var scope = factory.CreateSystemScope("Commit background job request")) { var dbContext = scope.ServiceProvider.GetRequiredService(); var jobs = scope.ServiceProvider.GetRequiredService(); await using var transaction = await dbContext.Database.BeginTransactionAsync(); committedJob = await jobs.EnqueueAsync(new CreateBackgroundJobCommand( tenantId, "tenant_domain_recheck", JsonSerializer.SerializeToElement(new { }))); await transaction.CommitAsync(); } BackgroundJobStatus? status = null; for (var attempt = 0; attempt < 60; attempt++) { using var verification = factory.CreateSystemScope("Wait for RabbitMQ background job consumer"); var dbContext = verification.ServiceProvider.GetRequiredService(); status = await dbContext.BackgroundJobs .Where(item => item.Id == committedJob.Id) .Select(item => (BackgroundJobStatus?)item.Status) .SingleAsync(); if (status == BackgroundJobStatus.Succeeded) break; await Task.Delay(250); } Assert.Equal(BackgroundJobStatus.Succeeded, status); using (var verification = factory.CreateSystemScope("Verify background job inbox and status view")) { var dbContext = verification.ServiceProvider.GetRequiredService(); Assert.Single(await dbContext.Set() .Where(item => item.MessageId == committedJob.Id) .ToArrayAsync()); var jobs = verification.ServiceProvider.GetRequiredService(); var statusView = await jobs.ListAsync(tenantId, "tenant_domain_recheck"); Assert.Contains(statusView, item => item.Id == committedJob.Id && item.Status == BackgroundJobStatus.Succeeded); } var managementEndpoint = ResolveRabbitManagementEndpoint(rabbitMqHost); if (managementEndpoint is not null) { Assert.Equal(0, await GetQueueMessageCountAsync( managementEndpoint, "background-job-requested_error")); } await worker.StopAsync(); } [Fact] public async Task Worker_consumer_uses_inbox_and_updates_non_authoritative_redis_version() { var rabbitMqHost = Environment.GetEnvironmentVariable("TIKU_TEST_RABBITMQ"); var redisConnection = Environment.GetEnvironmentVariable("TIKU_TEST_REDIS"); if (string.IsNullOrWhiteSpace(rabbitMqHost) || string.IsNullOrWhiteSpace(redisConnection)) { return; } await using var factory = CreateRabbitFactory(rabbitMqHost); using var client = factory.CreateClient(); Assert.True(await WaitForReadyAsync(client), "API dependencies did not become ready within 10 seconds."); var redisEnvironment = $"consumer-{Guid.NewGuid():N}"; var workerBuilder = Host.CreateApplicationBuilder(); workerBuilder.Services.AddApplication(); workerBuilder.Services.AddInfrastructure(factory.DatabaseConnectionString); workerBuilder.Services.AddRedisSecurity(redisConnection, redisEnvironment); workerBuilder.Services.AddReliableMessaging(CreateRabbitOptions(rabbitMqHost, configureConsumers: true)); using var worker = workerBuilder.Build(); await worker.StartAsync(); var tenantId = Guid.NewGuid(); var userId = Guid.NewGuid(); var messageId = Guid.NewGuid(); const long version = 123456789; using (var scope = factory.CreateSystemScope("Publish duplicate inbox test message")) { var dbContext = scope.ServiceProvider.GetRequiredService(); var publishEndpoint = scope.ServiceProvider.GetRequiredService(); await using var transaction = await dbContext.Database.BeginTransactionAsync(); var message = new AuthorizationStateChangedV1( Guid.NewGuid(), tenantId, userId, "consumer_test", version, DateTimeOffset.UtcNow, Guid.NewGuid().ToString("N")); await publishEndpoint.Publish(message, context => context.MessageId = messageId); await publishEndpoint.Publish(message, context => context.MessageId = messageId); await dbContext.SaveChangesAsync(); await transaction.CommitAsync(); } var redisKey = $"tiku:{redisEnvironment}:auth-inv:authorization:{tenantId:N}:{userId:N}"; var consumed = false; for (var attempt = 0; attempt < 60; attempt++) { var multiplexer = worker.Services.GetRequiredService(); if (await multiplexer.GetDatabase().StringGetAsync(redisKey) == version) { consumed = true; break; } await Task.Delay(250); } Assert.True(consumed, "Worker did not consume the security event within 15 seconds."); using (var verification = factory.CreateSystemScope("Verify duplicate consumer inbox")) { var dbContext = verification.ServiceProvider.GetRequiredService(); Assert.Single(await dbContext.Set() .Where(item => item.MessageId == messageId) .ToArrayAsync()); } var managementEndpoint = ResolveRabbitManagementEndpoint(rabbitMqHost); if (managementEndpoint is not null) { Assert.Equal(0, await GetQueueMessageCountAsync( managementEndpoint, "security-state-changed_error")); } var redis = worker.Services.GetRequiredService(); await redis.GetDatabase().KeyDeleteAsync(redisKey); await worker.StopAsync(); } [Fact] public async Task RabbitMq_health_and_bus_outbox_follow_database_transaction() { var rabbitMqHost = Environment.GetEnvironmentVariable("TIKU_TEST_RABBITMQ"); if (string.IsNullOrWhiteSpace(rabbitMqHost)) { return; } await using var factory = CreateRabbitFactory(rabbitMqHost); using var client = factory.CreateClient(); Assert.True(await WaitForReadyAsync(client), "API dependencies did not become ready within 10 seconds."); var ready = await client.GetAsync("/api/health/ready"); using var readyJson = JsonDocument.Parse(await ready.Content.ReadAsStringAsync()); Assert.Equal(HttpStatusCode.OK, ready.StatusCode); Assert.True(readyJson.RootElement.GetProperty("rabbitMq").GetProperty("configured").GetBoolean()); Assert.True(readyJson.RootElement.GetProperty("rabbitMq").GetProperty("ready").GetBoolean()); using (var rollbackScope = factory.CreateSystemScope("Verify rolled back bus outbox")) { var dbContext = rollbackScope.ServiceProvider.GetRequiredService(); var publisher = rollbackScope.ServiceProvider.GetRequiredService(); await using var transaction = await dbContext.Database.BeginTransactionAsync(); await publisher.AuthorizationChangedAsync( null, null, "rollback_test", 1, Guid.NewGuid().ToString("N")); await dbContext.SaveChangesAsync(); await transaction.RollbackAsync(); } using (var verification = factory.CreateSystemScope("Verify rolled back outbox is empty")) { var dbContext = verification.ServiceProvider.GetRequiredService(); Assert.Empty(await dbContext.Set().ToArrayAsync()); } using (var commitScope = factory.CreateSystemScope("Verify committed bus outbox")) { var dbContext = commitScope.ServiceProvider.GetRequiredService(); var publisher = commitScope.ServiceProvider.GetRequiredService(); await using var transaction = await dbContext.Database.BeginTransactionAsync(); await publisher.AuthorizationChangedAsync( null, null, "commit_test", 2, Guid.NewGuid().ToString("N")); await dbContext.SaveChangesAsync(); await transaction.CommitAsync(); } var drained = false; for (var attempt = 0; attempt < 40; attempt++) { using var verification = factory.CreateSystemScope("Wait for committed outbox delivery"); var dbContext = verification.ServiceProvider.GetRequiredService(); if (!await dbContext.Set().AnyAsync()) { drained = true; break; } await Task.Delay(250); } if (!drained) { using var diagnostics = factory.CreateSystemScope("Inspect undelivered outbox"); var dbContext = diagnostics.ServiceProvider.GetRequiredService(); var messages = await dbContext.Set().CountAsync(); var states = await dbContext.Set().CountAsync(); Assert.Fail($"Committed MassTransit outbox was not delivered within 10 seconds. messages={messages}, states={states}"); } } private static Api.ApiTestFactory CreateRabbitFactory(string rabbitMqHost) => new(configurationOverrides: new Dictionary { ["RabbitMq:Host"] = rabbitMqHost, ["RabbitMq:Username"] = Environment.GetEnvironmentVariable("TIKU_TEST_RABBITMQ_USERNAME") ?? "guest", ["RabbitMq:Password"] = Environment.GetEnvironmentVariable("TIKU_TEST_RABBITMQ_PASSWORD") ?? "guest" }); private static MessagingOptions CreateRabbitOptions(string rabbitMqHost, bool configureConsumers) => new() { Host = rabbitMqHost, Username = Environment.GetEnvironmentVariable("TIKU_TEST_RABBITMQ_USERNAME") ?? "guest", Password = Environment.GetEnvironmentVariable("TIKU_TEST_RABBITMQ_PASSWORD") ?? "guest", ConfigureConsumers = configureConsumers }; private static async Task WaitForReadyAsync(HttpClient client, int attempts = 40) { for (var attempt = 0; attempt < attempts; attempt++) { if ((await client.GetAsync("/api/health/ready")).StatusCode == HttpStatusCode.OK) { return true; } await Task.Delay(250); } return false; } private static Uri? ResolveRabbitManagementEndpoint(string rabbitMqHost) { var configured = Environment.GetEnvironmentVariable("TIKU_TEST_RABBITMQ_MANAGEMENT"); if (!string.IsNullOrWhiteSpace(configured)) { return new Uri(configured); } var broker = new Uri(rabbitMqHost); return broker.IsLoopback ? new Uri($"http://{broker.Host}:15672") : null; } private static async Task GetQueueMessageCountAsync(Uri managementEndpoint, string queueName) { using var client = new HttpClient { BaseAddress = managementEndpoint }; var username = Environment.GetEnvironmentVariable("TIKU_TEST_RABBITMQ_USERNAME") ?? "guest"; var password = Environment.GetEnvironmentVariable("TIKU_TEST_RABBITMQ_PASSWORD") ?? "guest"; var credentials = Convert.ToBase64String(System.Text.Encoding.UTF8.GetBytes($"{username}:{password}")); client.DefaultRequestHeaders.Authorization = new AuthenticationHeaderValue("Basic", credentials); using var response = await client.GetAsync($"/api/queues/%2F/{Uri.EscapeDataString(queueName)}"); if (response.StatusCode == HttpStatusCode.NotFound) { return 0; } response.EnsureSuccessStatusCode(); using var document = JsonDocument.Parse(await response.Content.ReadAsStringAsync()); return document.RootElement.GetProperty("messages").GetInt32(); } }