forked from gongxuegit/tiku-backend.net
feat(security): complete capability messaging workflows
This commit is contained in:
@@ -4,11 +4,16 @@ 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;
|
||||
@@ -17,6 +22,199 @@ 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<TikuDbContext>();
|
||||
var publisher = scope.ServiceProvider.GetRequiredService<ISecurityEventPublisher>();
|
||||
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<TikuDbContext>();
|
||||
Assert.NotEmpty(await verificationDbContext.Set<OutboxMessage>().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<TikuDbContext>();
|
||||
if (!await dbContext.Set<OutboxMessage>().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<bool> 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<TikuDbContext>();
|
||||
var jobs = scope.ServiceProvider.GetRequiredService<IBackgroundJobService>();
|
||||
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<TikuDbContext>();
|
||||
Assert.False(await dbContext.BackgroundJobs.AnyAsync(item => item.Id == rolledBackJobId));
|
||||
Assert.False(await dbContext.Set<OutboxMessage>().AnyAsync(
|
||||
item => item.MessageId == rolledBackJobId));
|
||||
}
|
||||
|
||||
BackgroundJobItem committedJob;
|
||||
using (var scope = factory.CreateSystemScope("Commit background job request"))
|
||||
{
|
||||
var dbContext = scope.ServiceProvider.GetRequiredService<TikuDbContext>();
|
||||
var jobs = scope.ServiceProvider.GetRequiredService<IBackgroundJobService>();
|
||||
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<TikuDbContext>();
|
||||
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<TikuDbContext>();
|
||||
Assert.Single(await dbContext.Set<InboxState>()
|
||||
.Where(item => item.MessageId == committedJob.Id)
|
||||
.ToArrayAsync());
|
||||
var jobs = verification.ServiceProvider.GetRequiredService<IBackgroundJobService>();
|
||||
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()
|
||||
{
|
||||
@@ -177,9 +375,9 @@ public sealed class MassTransitOutboxTests
|
||||
ConfigureConsumers = configureConsumers
|
||||
};
|
||||
|
||||
private static async Task<bool> WaitForReadyAsync(HttpClient client)
|
||||
private static async Task<bool> WaitForReadyAsync(HttpClient client, int attempts = 40)
|
||||
{
|
||||
for (var attempt = 0; attempt < 40; attempt++)
|
||||
for (var attempt = 0; attempt < attempts; attempt++)
|
||||
{
|
||||
if ((await client.GetAsync("/api/health/ready")).StatusCode == HttpStatusCode.OK)
|
||||
{
|
||||
@@ -210,6 +408,10 @@ public sealed class MassTransitOutboxTests
|
||||
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();
|
||||
|
||||
Reference in New Issue
Block a user