/
afanasevn
/
MaxSystems
Обзор
Документация
Войти
/
afanasevn
/
MaxSystems
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
tests/MaxSystems.IntegrationTests/IntegrationEventBusTests.cs
353 строки
15 KB
IBS\NAfanasev
Update solution references
23 июн 2026, 16:24
23 июн 2026, 16:24
02b1783
Код
Авторство
О чём код?
using System.Net; using System.Net.Http.Json; using System.Text.Json; using MaxSystems.Application.Integration; using MaxSystems.Application.Integration.Events; using MaxSystems.Domain.Enums; using MaxSystems.IntegrationTests.Infrastructure; namespace MaxSystems.IntegrationTests; /// <summary> /// Интеграционные тесты publisher-контура Integration Event Bus: /// API => transactional outbox (Wolverine) => RabbitMQ + журнал integration_delivery_logs. /// </summary> [Collection(IntegrationPublisherTestCollection.Name)] public sealed class IntegrationEventBusTests(IntegrationPublisherTestFixture fixture) { private static readonly JsonSerializerOptions JsonOptions = new() { PropertyNameCaseInsensitive = true, }; /// <summary> /// Создание работы с исполнителем порождает два типа событий, каждый — в две целевые очереди. /// </summary> /// <remarks> /// Проверяется: /// <list type="bullet"> /// <item>4 записи в integration_delivery_logs со статусом <see cref="IntegrationDeliveryStatus.Pending"/></item> /// <item>work.created и work.assignee.changed — по одной строке на SD и Notifications</item> /// <item>корректные routing keys и имена очередей из <see cref="IntegrationTopology"/></item> /// </list> /// RabbitMQ в этом тесте не проверяется (см. <see cref="CreateWork_OutboxDispatchesToRabbitMqExchange"/>). /// </remarks> [Fact] public async Task CreateWork_WithAssignee_WritesFourPendingDeliveryLogRows() { await fixture.LoginAsManagerAsync(); // AssigneeId задан => WorkService публикует work.created и work.assignee.changed. var response = await fixture.Client.PostAsJsonAsync( "/api/works", new { title = "Integration bus create test", type = WorkType.Repair, category = "Integration test", system = EquipmentSystem.Ups, priority = Priority.P3, plannedDate = DateTime.UtcNow.Date.AddDays(3), assigneeId = TestSeedIds.Spec1, }, JsonOptions); Assert.Equal(HttpStatusCode.OK, response.StatusCode); using var created = await response.Content.ReadFromJsonAsync<JsonDocument>(); var workId = created!.RootElement.GetProperty("id").GetGuid(); var logs = await IntegrationTestDatabase.GetDeliveryLogsForWorkAsync( fixture.PostgresConnectionString, workId); // 2 event types × 2 queues (ServiceDesk + Notifications). Assert.Equal(4, logs.Count); Assert.Equal(2, logs.Count(x => x.EventType == IntegrationEventTypes.WorkCreated)); Assert.Equal(2, logs.Count(x => x.EventType == IntegrationEventTypes.WorkAssigneeChanged)); Assert.All(logs, row => Assert.Equal(IntegrationDeliveryStatus.Pending, row.Status)); foreach (var eventType in new[] { IntegrationEventTypes.WorkCreated, IntegrationEventTypes.WorkAssigneeChanged, }) { var queues = logs .Where(x => x.EventType == eventType) .Select(x => x.QueueName) .ToHashSet(); Assert.Contains(IntegrationTopology.ServiceDeskQueue, queues); Assert.Contains(IntegrationTopology.NotificationsQueue, queues); Assert.Equal(2, queues.Count); } Assert.All( logs.Where(x => x.EventType == IntegrationEventTypes.WorkCreated), row => Assert.Equal(IntegrationRoutingKeys.WorkCreated, row.RoutingKey)); Assert.All( logs.Where(x => x.EventType == IntegrationEventTypes.WorkAssigneeChanged), row => Assert.Equal(IntegrationRoutingKeys.WorkAssigneeChanged, row.RoutingKey)); } /// <summary> /// PATCH статуса работы пишет journal и публикует событие с динамическим routing key. /// </summary> /// <remarks> /// Проверяется: /// <list type="bullet"> /// <item>2 строки journal для work.status.inprogress</item> /// <item>сообщение в exchange <see cref="IntegrationTopology.ExchangeName"/> с ожидаемым routing key</item> /// <item>payload <see cref="WorkStatusChangedIntegrationEvent"/> (PreviousStatus = New, NewStatus = InProgress)</item> /// </list> /// Используется seed-работа <see cref="TestSeedIds.Work2"/> (статус New). /// </remarks> [Fact] public async Task UpdateWorkStatus_WritesJournalWithDynamicRoutingKey() { await fixture.LoginAsManagerAsync(); const WorkStatus newStatus = WorkStatus.InProgress; var expectedRoutingKey = IntegrationRoutingKeys.WorkStatusChanged(newStatus); // Подписываемся на exchange до PATCH, чтобы не пропустить dispatch. await using var capture = await RabbitMqMessageCapture.BindToExchangeAsync( fixture.RabbitMqAmqpUri, IntegrationTopology.ExchangeName, routingKeyPattern: expectedRoutingKey); var response = await fixture.Client.PatchAsJsonAsync( $"/api/works/{TestSeedIds.Work2}", new { status = newStatus }, JsonOptions); Assert.Equal(HttpStatusCode.OK, response.StatusCode); var logs = await IntegrationTestDatabase.GetDeliveryLogsForWorkAsync( fixture.PostgresConnectionString, TestSeedIds.Work2); var statusLogs = logs .Where(x => x.EventType == IntegrationEventTypes.WorkStatusChanged) .ToList(); Assert.Equal(2, statusLogs.Count); Assert.All(statusLogs, row => Assert.Equal(IntegrationDeliveryStatus.Pending, row.Status)); Assert.All(statusLogs, row => Assert.Equal(expectedRoutingKey, row.RoutingKey)); var queues = statusLogs.Select(x => x.QueueName).ToHashSet(); Assert.Contains(IntegrationTopology.ServiceDeskQueue, queues); Assert.Contains(IntegrationTopology.NotificationsQueue, queues); Assert.Equal(2, queues.Count); var statusMessage = await capture.WaitForMessageAsync( message => message.RoutingKey == expectedRoutingKey, timeout: TimeSpan.FromSeconds(30)); using var payload = RabbitMqMessageCapture.ParseJson(statusMessage.Body); var evt = payload.Deserialize<WorkStatusChangedIntegrationEvent>(JsonOptions); Assert.NotNull(evt); Assert.Equal(TestSeedIds.Work2, evt!.WorkId); Assert.Equal(WorkStatus.New, evt.PreviousStatus); Assert.Equal(newStatus, evt.NewStatus); } /// <summary> /// Завершение чек-листа переводит работу в OnReview и публикует work.status.onreview. /// </summary> /// <remarks> /// Проверяется путь SaveChecklistAsync => IntegrationEventPublisher: /// journal (2 строки) + RabbitMQ payload с PreviousStatus = InProgress. /// Seed: <see cref="TestSeedIds.Work1"/> (InProgress, чек-лист не завершён). /// </remarks> [Fact] public async Task CompleteChecklist_WritesJournalWithOnReviewRoutingKey() { await fixture.LoginAsManagerAsync(); const WorkStatus expectedNewStatus = WorkStatus.OnReview; var expectedRoutingKey = IntegrationRoutingKeys.WorkStatusChanged(expectedNewStatus); await using var capture = await RabbitMqMessageCapture.BindToExchangeAsync( fixture.RabbitMqAmqpUri, IntegrationTopology.ExchangeName, routingKeyPattern: expectedRoutingKey); var response = await fixture.Client.PostAsJsonAsync( $"/api/works/{TestSeedIds.Work1}/checklist", new { type = "UPS_MAINTENANCE", itemsJson = """{"1":"yes","4":"42","6":"Норма"}""", summaryJson = """{"battery":"Норма","needZip":"Нет"}""", completed = true, }, JsonOptions); Assert.Equal(HttpStatusCode.OK, response.StatusCode); var logs = await IntegrationTestDatabase.GetDeliveryLogsForWorkAsync( fixture.PostgresConnectionString, TestSeedIds.Work1); var statusLogs = logs .Where(x => x.EventType == IntegrationEventTypes.WorkStatusChanged && x.RoutingKey == expectedRoutingKey) .ToList(); Assert.Equal(2, statusLogs.Count); Assert.All(statusLogs, row => Assert.Equal(IntegrationDeliveryStatus.Pending, row.Status)); var message = await capture.WaitForMessageAsync( m => m.RoutingKey == expectedRoutingKey, timeout: TimeSpan.FromSeconds(30)); using var payload = RabbitMqMessageCapture.ParseJson(message.Body); var evt = payload.Deserialize<WorkStatusChangedIntegrationEvent>(JsonOptions); Assert.NotNull(evt); Assert.Equal(TestSeedIds.Work1, evt!.WorkId); Assert.Equal(WorkStatus.InProgress, evt.PreviousStatus); Assert.Equal(expectedNewStatus, evt.NewStatus); } /// <summary> /// Transactional outbox доставляет work.created в topic exchange maxsystems.events. /// </summary> /// <remarks> /// Проверяется: /// <list type="bullet"> /// <item>сообщение с routing key <see cref="IntegrationRoutingKeys.WorkCreated"/> получено тестовой очередью</item> /// <item>outbox либо кратковременно растёт, либо успевает опустошиться после dispatch</item> /// <item>тело содержит идентификатор работы</item> /// </list> /// </remarks> [Fact] public async Task CreateWork_OutboxDispatchesToRabbitMqExchange() { await fixture.LoginAsManagerAsync(); var outboxBefore = await IntegrationTestDatabase.CountOutgoingEnvelopesAsync(fixture.PostgresConnectionString); await using var capture = await RabbitMqMessageCapture.BindToExchangeAsync( fixture.RabbitMqAmqpUri, IntegrationTopology.ExchangeName, routingKeyPattern: "work.#"); // Без assignee — только одно событие work.created. var response = await fixture.Client.PostAsJsonAsync( "/api/works", new { title = "Outbox dispatch test", type = WorkType.Ticket, category = "Integration test", system = EquipmentSystem.Acs, priority = Priority.P2, plannedDate = DateTime.UtcNow.Date.AddDays(1), }, JsonOptions); Assert.Equal(HttpStatusCode.OK, response.StatusCode); var message = await capture.WaitForMessageAsync( m => m.RoutingKey == IntegrationRoutingKeys.WorkCreated, timeout: TimeSpan.FromSeconds(30)); // Dispatch может завершиться быстрее, чем poll увидит строку в outbox. var sawOutboxEntry = await WaitForOutboxGrowthAsync(outboxBefore, TimeSpan.FromSeconds(1)); var outboxAfter = await IntegrationTestDatabase.CountOutgoingEnvelopesAsync(fixture.PostgresConnectionString); Assert.True( sawOutboxEntry || outboxAfter <= outboxBefore, $"Expected outbox growth or drain after dispatch (before={outboxBefore}, after={outboxAfter})."); Assert.Equal(IntegrationRoutingKeys.WorkCreated, message.RoutingKey); Assert.StartsWith("work.", message.RoutingKey, StringComparison.Ordinal); using var payload = RabbitMqMessageCapture.ParseJson(message.Body); Assert.True(payload.RootElement.TryGetProperty("workId", out _) || payload.RootElement.TryGetProperty("WorkId", out _)); } /// <summary> /// При ошибке SaveChanges (FK) journal и outbox откатываются вместе с транзакцией EF. /// </summary> /// <remarks> /// Проверяется: /// <list type="bullet"> /// <item>запрос с несуществующим AssigneeId завершается ошибкой (HTTP 500 или исключение TestServer)</item> /// <item>число строк в integration_delivery_logs и wolverine.wolverine_outgoing_envelopes не меняется</item> /// </list> /// </remarks> [Fact] public async Task CreateWork_WhenSaveChangesFails_RollsBackJournalAndOutbox() { await fixture.LoginAsManagerAsync(); var logsBefore = await IntegrationTestDatabase.CountDeliveryLogsAsync(fixture.PostgresConnectionString); var outboxBefore = await IntegrationTestDatabase.CountOutgoingEnvelopesAsync(fixture.PostgresConnectionString); // Пользователя с таким Id нет => FK works.AssigneeId => users. var invalidAssigneeId = Guid.Parse("aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa"); HttpResponseMessage? response = null; Exception? failure = null; try { response = await fixture.Client.PostAsJsonAsync( "/api/works", new { title = "Should fail FK", type = WorkType.Repair, category = "Integration test", system = EquipmentSystem.Ups, priority = Priority.P3, plannedDate = DateTime.UtcNow.Date.AddDays(2), assigneeId = invalidAssigneeId, }, JsonOptions); } catch (Exception ex) { // WebApplicationFactory может пробросить DbUpdateException на клиент. failure = ex; } Assert.True( failure is not null || response is { IsSuccessStatusCode: false }, "Expected FK violation to surface as HTTP error or thrown exception."); if (response is not null) Assert.Equal(HttpStatusCode.InternalServerError, response.StatusCode); var logsAfter = await IntegrationTestDatabase.CountDeliveryLogsAsync(fixture.PostgresConnectionString); var outboxAfter = await IntegrationTestDatabase.CountOutgoingEnvelopesAsync(fixture.PostgresConnectionString); Assert.Equal(logsBefore, logsAfter); Assert.Equal(outboxBefore, outboxAfter); } /// <summary> /// Краткий poll таблицы outbox — dispatch асинхронный, окно может не поймать строку. /// </summary> private async Task<bool> WaitForOutboxGrowthAsync(int baseline, TimeSpan timeout) { var deadline = DateTime.UtcNow + timeout; while (DateTime.UtcNow < deadline) { var count = await IntegrationTestDatabase.CountOutgoingEnvelopesAsync(fixture.PostgresConnectionString); if (count > baseline) return true; await Task.Delay(TimeSpan.FromMilliseconds(50)); } return false; } }