/
ivanstrike
/
tasker
Обзор
Документация
Войти
/
ivanstrike
/
tasker
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
Lab4/Server/Services/MessageProcessingService.cs
700 строк
26 KB
ivanstrike
add DLQ simulation
23 дек 2025, 14:22
23 дек 2025, 14:22
e6b7124
Код
Авторство
О чём код?
using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using Newtonsoft.Json; using Newtonsoft.Json.Linq; using RabbitMQ.Client; using RabbitMQ.Client.Events; using System.Text; using TaskerMQ.Server.Configuration; using TaskerMQ.Server.Data; using TaskerMQ.Server.DTOs; using TaskerMQ.Server.Models; namespace TaskerMQ.Server.Services; public class MessageProcessingService : BackgroundService { private readonly IServiceProvider _serviceProvider; private readonly RabbitMQSettings _rabbitSettings; private readonly ILogger<MessageProcessingService> _logger; private readonly RabbitMQConnectionFactory _connectionFactory; public MessageProcessingService( IServiceProvider serviceProvider, RabbitMQSettings rabbitSettings, ILogger<MessageProcessingService> logger, RabbitMQConnectionFactory connectionFactory) { _serviceProvider = serviceProvider; _rabbitSettings = rabbitSettings; _logger = logger; _connectionFactory = connectionFactory; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { _logger.LogInformation("Message Processing Service started"); try { var channel = await _connectionFactory.GetChannelAsync(); // Declare queues await channel.QueueDeclareAsync( queue: _rabbitSettings.RequestQueue, durable: true, exclusive: false, autoDelete: false, arguments: new Dictionary<string, object?> { { "x-dead-letter-exchange", "" }, { "x-dead-letter-routing-key", _rabbitSettings.DeadLetterQueue } }); await channel.QueueDeclareAsync( queue: _rabbitSettings.ResponseQueue, durable: true, exclusive: false, autoDelete: false, arguments: null); await channel.QueueDeclareAsync( queue: _rabbitSettings.DeadLetterQueue, durable: true, exclusive: false, autoDelete: false, arguments: null); var consumer = new AsyncEventingBasicConsumer(channel); consumer.ReceivedAsync += async (model, ea) => { await ProcessMessageAsync(ea, channel); }; await channel.BasicConsumeAsync( queue: _rabbitSettings.RequestQueue, autoAck: false, consumer: consumer); _logger.LogInformation( "Listening for messages on queue: {Queue}", _rabbitSettings.RequestQueue); await Task.Delay(Timeout.Infinite, stoppingToken); } catch (Exception ex) { _logger.LogError(ex, "Error in Message Processing Service"); throw; } } private async Task ProcessMessageAsync(BasicDeliverEventArgs ea, IChannel channel) { var body = ea.Body.ToArray(); var message = Encoding.UTF8.GetString(body); _logger.LogInformation("Received message: {Message}", message); try { var request = JsonConvert.DeserializeObject<RequestMessage>(message); if (request == null) { _logger.LogWarning("Failed to deserialize message"); await channel.BasicRejectAsync(ea.DeliveryTag, false); return; } using var scope = _serviceProvider.CreateScope(); var authService = scope.ServiceProvider.GetRequiredService<AuthenticationService>(); var idempotencyService = scope.ServiceProvider.GetRequiredService<IdempotencyService>(); // Check authentication if (!authService.ValidateApiKey(request.Auth)) { var errorResponse = new ResponseMessage { CorrelationId = request.Id, Status = "error", Data = null, Error = "Unauthorized: Invalid API key" }; await SendResponseAsync(errorResponse, channel); await channel.BasicAckAsync(ea.DeliveryTag, false); return; } // Check idempotency var cachedResponse = await idempotencyService.GetCachedResponseAsync(request.Id); if (cachedResponse != null) { _logger.LogInformation( "Returning cached response for request {RequestId}", request.Id); await SendResponseAsync(cachedResponse, channel); await channel.BasicAckAsync(ea.DeliveryTag, false); return; } // Process request var response = await ProcessRequestAsync(request, scope); // Cache response var responseJson = JsonConvert.SerializeObject(response); await idempotencyService.CacheResponseAsync(request.Id, responseJson); // Send response await SendResponseAsync(response, channel); await channel.BasicAckAsync(ea.DeliveryTag, false); _logger.LogInformation( "Successfully processed request {RequestId}", request.Id); } catch (Exception ex) { _logger.LogError(ex, "Error processing message"); // Check retry count var retryCount = GetRetryCount(ea); if (retryCount < 2) // Will retry 2 times (total 3 attempts) { _logger.LogInformation( "Retrying message (attempt {Retry} of 3)", retryCount + 2); // Increment retry count and republish var properties = new BasicProperties { Headers = new Dictionary<string, object?> { { "x-retry-count", retryCount + 1 } } }; await channel.BasicPublishAsync( exchange: "", routingKey: _rabbitSettings.RequestQueue, mandatory: false, basicProperties: properties, body: ea.Body); await channel.BasicAckAsync(ea.DeliveryTag, false); } else { _logger.LogWarning("Max retries (3) reached, sending to DLQ"); await channel.BasicRejectAsync(ea.DeliveryTag, false); } } } private async Task<ResponseMessage> ProcessRequestAsync( RequestMessage request, IServiceScope scope) { var context = scope.ServiceProvider.GetRequiredService<ApplicationDbContext>(); // Special handling for DLQ demonstration - BEFORE try-catch if (request.Action == "demonstrate_dlq") { return new ResponseMessage { CorrelationId = request.Id, Status = "ok", Data = await DemonstrateDLQAsync(request, scope), Error = null }; } // Simulate failure for DLQ testing - throw exception to trigger retry mechanism // This MUST be before try-catch to properly trigger retry/DLQ logic if (request.Action == "simulate_failure") { _logger.LogWarning("Simulating failure for DLQ demonstration - throwing exception"); throw new InvalidOperationException("Simulated failure for DLQ testing"); } try { object? result = null; var actionParts = request.Action.Split('_'); var operation = actionParts[0]; // create, get, update, delete, list var entity = actionParts.Length > 1 ? actionParts[1] : ""; switch (entity) { case "user": result = await ProcessUserActionAsync(operation, request, context); break; case "project": result = await ProcessProjectActionAsync(operation, request, context); break; case "task": result = await ProcessTaskActionAsync(operation, request, context); break; case "users": result = await ListUsersAsync(context); break; case "projects": result = await ListProjectsAsync(request, context); break; case "tasks": result = await ListTasksAsync(request, context); break; default: throw new InvalidOperationException($"Unknown action: {request.Action}"); } return new ResponseMessage { CorrelationId = request.Id, Status = "ok", Data = result, Error = null }; } catch (Exception ex) { _logger.LogError(ex, "Error processing action {Action}", request.Action); return new ResponseMessage { CorrelationId = request.Id, Status = "error", Data = null, Error = ex.Message }; } } private async Task<object> ProcessUserActionAsync( string operation, RequestMessage request, ApplicationDbContext context) { switch (operation) { case "create": var createDto = JsonConvert.DeserializeObject<CreateUserDto>( request.Data?.ToString() ?? "{}"); if (createDto == null) throw new InvalidOperationException("Invalid user data"); var user = new User { Username = createDto.Username, Email = createDto.Email, PasswordHash = BCrypt.Net.BCrypt.HashPassword(createDto.Password), CreatedAt = DateTime.UtcNow, UpdatedAt = DateTime.UtcNow }; context.Users.Add(user); await context.SaveChangesAsync(); return new UserResponseDto { Id = user.Id, Username = user.Username, Email = user.Email, CreatedAt = user.CreatedAt, UpdatedAt = user.UpdatedAt }; case "get": var userId = GetIdFromData(request.Data); var existingUser = await context.Users.FindAsync(userId); if (existingUser == null) throw new InvalidOperationException($"User {userId} not found"); return new UserResponseDto { Id = existingUser.Id, Username = existingUser.Username, Email = existingUser.Email, CreatedAt = existingUser.CreatedAt, UpdatedAt = existingUser.UpdatedAt }; case "update": userId = GetIdFromData(request.Data); existingUser = await context.Users.FindAsync(userId); if (existingUser == null) throw new InvalidOperationException($"User {userId} not found"); var updateDto = JsonConvert.DeserializeObject<UpdateUserDto>( request.Data?.ToString() ?? "{}"); if (updateDto == null) throw new InvalidOperationException("Invalid update data"); if (!string.IsNullOrEmpty(updateDto.Username)) existingUser.Username = updateDto.Username; if (!string.IsNullOrEmpty(updateDto.Email)) existingUser.Email = updateDto.Email; if (!string.IsNullOrEmpty(updateDto.Password)) existingUser.PasswordHash = BCrypt.Net.BCrypt.HashPassword(updateDto.Password); existingUser.UpdatedAt = DateTime.UtcNow; await context.SaveChangesAsync(); return new UserResponseDto { Id = existingUser.Id, Username = existingUser.Username, Email = existingUser.Email, CreatedAt = existingUser.CreatedAt, UpdatedAt = existingUser.UpdatedAt }; case "delete": userId = GetIdFromData(request.Data); existingUser = await context.Users.FindAsync(userId); if (existingUser == null) throw new InvalidOperationException($"User {userId} not found"); context.Users.Remove(existingUser); await context.SaveChangesAsync(); return new { message = "User deleted successfully" }; default: throw new InvalidOperationException($"Unknown operation: {operation}"); } } private async Task<object> ProcessProjectActionAsync( string operation, RequestMessage request, ApplicationDbContext context) { switch (operation) { case "create": var createDto = JsonConvert.DeserializeObject<CreateProjectDto>( request.Data?.ToString() ?? "{}"); if (createDto == null) throw new InvalidOperationException("Invalid project data"); var project = new Project { Name = createDto.Name, Description = createDto.Description, UserId = createDto.UserId, CreatedAt = DateTime.UtcNow, UpdatedAt = DateTime.UtcNow }; context.Projects.Add(project); await context.SaveChangesAsync(); return new ProjectResponseDto { Id = project.Id, Name = project.Name, Description = project.Description, UserId = project.UserId, CreatedAt = project.CreatedAt, UpdatedAt = project.UpdatedAt }; case "get": var projectId = GetIdFromData(request.Data); var existingProject = await context.Projects.FindAsync(projectId); if (existingProject == null) throw new InvalidOperationException($"Project {projectId} not found"); return new ProjectResponseDto { Id = existingProject.Id, Name = existingProject.Name, Description = existingProject.Description, UserId = existingProject.UserId, CreatedAt = existingProject.CreatedAt, UpdatedAt = existingProject.UpdatedAt }; case "update": projectId = GetIdFromData(request.Data); existingProject = await context.Projects.FindAsync(projectId); if (existingProject == null) throw new InvalidOperationException($"Project {projectId} not found"); var updateDto = JsonConvert.DeserializeObject<UpdateProjectDto>( request.Data?.ToString() ?? "{}"); if (updateDto == null) throw new InvalidOperationException("Invalid update data"); if (!string.IsNullOrEmpty(updateDto.Name)) existingProject.Name = updateDto.Name; if (updateDto.Description != null) existingProject.Description = updateDto.Description; existingProject.UpdatedAt = DateTime.UtcNow; await context.SaveChangesAsync(); return new ProjectResponseDto { Id = existingProject.Id, Name = existingProject.Name, Description = existingProject.Description, UserId = existingProject.UserId, CreatedAt = existingProject.CreatedAt, UpdatedAt = existingProject.UpdatedAt }; case "delete": projectId = GetIdFromData(request.Data); existingProject = await context.Projects.FindAsync(projectId); if (existingProject == null) throw new InvalidOperationException($"Project {projectId} not found"); context.Projects.Remove(existingProject); await context.SaveChangesAsync(); return new { message = "Project deleted successfully" }; default: throw new InvalidOperationException($"Unknown operation: {operation}"); } } private async Task<object> ProcessTaskActionAsync( string operation, RequestMessage request, ApplicationDbContext context) { switch (operation) { case "create": var createDto = JsonConvert.DeserializeObject<CreateTaskDto>( request.Data?.ToString() ?? "{}"); if (createDto == null) throw new InvalidOperationException("Invalid task data"); var task = new TaskItem { Title = createDto.Title, Description = createDto.Description, ProjectId = createDto.ProjectId, UserId = createDto.UserId, DueDate = createDto.DueDate, Status = Models.TaskStatus.Pending, CreatedAt = DateTime.UtcNow, UpdatedAt = DateTime.UtcNow }; context.Tasks.Add(task); await context.SaveChangesAsync(); return new TaskResponseDto { Id = task.Id, Title = task.Title, Description = task.Description, Status = task.Status.ToString(), ProjectId = task.ProjectId, UserId = task.UserId, CreatedAt = task.CreatedAt, UpdatedAt = task.UpdatedAt, DueDate = task.DueDate }; case "get": var taskId = GetIdFromData(request.Data); var existingTask = await context.Tasks.FindAsync(taskId); if (existingTask == null) throw new InvalidOperationException($"Task {taskId} not found"); return new TaskResponseDto { Id = existingTask.Id, Title = existingTask.Title, Description = existingTask.Description, Status = existingTask.Status.ToString(), ProjectId = existingTask.ProjectId, UserId = existingTask.UserId, CreatedAt = existingTask.CreatedAt, UpdatedAt = existingTask.UpdatedAt, DueDate = existingTask.DueDate }; case "update": taskId = GetIdFromData(request.Data); existingTask = await context.Tasks.FindAsync(taskId); if (existingTask == null) throw new InvalidOperationException($"Task {taskId} not found"); var updateDto = JsonConvert.DeserializeObject<UpdateTaskDto>( request.Data?.ToString() ?? "{}"); if (updateDto == null) throw new InvalidOperationException("Invalid update data"); if (!string.IsNullOrEmpty(updateDto.Title)) existingTask.Title = updateDto.Title; if (updateDto.Description != null) existingTask.Description = updateDto.Description; if (!string.IsNullOrEmpty(updateDto.Status) && Enum.TryParse<Models.TaskStatus>(updateDto.Status, out var status)) existingTask.Status = status; if (updateDto.DueDate.HasValue) existingTask.DueDate = updateDto.DueDate; existingTask.UpdatedAt = DateTime.UtcNow; await context.SaveChangesAsync(); return new TaskResponseDto { Id = existingTask.Id, Title = existingTask.Title, Description = existingTask.Description, Status = existingTask.Status.ToString(), ProjectId = existingTask.ProjectId, UserId = existingTask.UserId, CreatedAt = existingTask.CreatedAt, UpdatedAt = existingTask.UpdatedAt, DueDate = existingTask.DueDate }; case "delete": taskId = GetIdFromData(request.Data); existingTask = await context.Tasks.FindAsync(taskId); if (existingTask == null) throw new InvalidOperationException($"Task {taskId} not found"); context.Tasks.Remove(existingTask); await context.SaveChangesAsync(); return new { message = "Task deleted successfully" }; default: throw new InvalidOperationException($"Unknown operation: {operation}"); } } private async Task<object> ListUsersAsync(ApplicationDbContext context) { var users = await context.Users.ToListAsync(); return users.Select(u => new UserResponseDto { Id = u.Id, Username = u.Username, Email = u.Email, CreatedAt = u.CreatedAt, UpdatedAt = u.UpdatedAt }).ToList(); } private async Task<object> ListProjectsAsync(RequestMessage request, ApplicationDbContext context) { IQueryable<Project> query = context.Projects; if (request.Data != null) { var jObj = JObject.Parse(request.Data.ToString() ?? "{}"); if (jObj["user_id"] != null) { var userId = jObj["user_id"]!.Value<int>(); query = query.Where(p => p.UserId == userId); } } var projects = await query.ToListAsync(); return projects.Select(p => new ProjectResponseDto { Id = p.Id, Name = p.Name, Description = p.Description, UserId = p.UserId, CreatedAt = p.CreatedAt, UpdatedAt = p.UpdatedAt }).ToList(); } private async Task<object> ListTasksAsync(RequestMessage request, ApplicationDbContext context) { IQueryable<TaskItem> query = context.Tasks; if (request.Data != null) { var jObj = JObject.Parse(request.Data.ToString() ?? "{}"); if (jObj["project_id"] != null) { var projectId = jObj["project_id"]!.Value<int>(); query = query.Where(t => t.ProjectId == projectId); } if (jObj["user_id"] != null) { var userId = jObj["user_id"]!.Value<int>(); query = query.Where(t => t.UserId == userId); } if (jObj["status"] != null && Enum.TryParse<Models.TaskStatus>(jObj["status"]!.Value<string>(), out var status)) { query = query.Where(t => t.Status == status); } } var tasks = await query.ToListAsync(); return tasks.Select(t => new TaskResponseDto { Id = t.Id, Title = t.Title, Description = t.Description, Status = t.Status.ToString(), ProjectId = t.ProjectId, UserId = t.UserId, CreatedAt = t.CreatedAt, UpdatedAt = t.UpdatedAt, DueDate = t.DueDate }).ToList(); } private int GetIdFromData(object? data) { if (data == null) throw new InvalidOperationException("Missing ID in request data"); var jObj = JObject.Parse(data.ToString() ?? "{}"); return jObj["id"]?.Value<int>() ?? throw new InvalidOperationException("Missing ID in request data"); } private async Task<object> DemonstrateDLQAsync(RequestMessage request, IServiceScope scope) { var channel = await _connectionFactory.GetChannelAsync(); // Получаем количество сообщений в DLQ var dlqInfo = await channel.QueueDeclarePassiveAsync(_rabbitSettings.DeadLetterQueue); var messageCount = dlqInfo.MessageCount; _logger.LogInformation("DLQ contains {MessageCount} messages", messageCount); return new { dlq_queue = _rabbitSettings.DeadLetterQueue, message_count = messageCount, description = "Dead Letter Queue is used for messages that failed processing after maximum retry attempts", info = new { max_retries = 3, retry_mechanism = "automatic requeue on failure", dlq_purpose = "Store failed messages for manual inspection and reprocessing" } }; } private async Task SendResponseAsync(object response, IChannel channel) { var responseJson = response is string str ? str : JsonConvert.SerializeObject(response); var responseBytes = Encoding.UTF8.GetBytes(responseJson); await channel.BasicPublishAsync( exchange: "", routingKey: _rabbitSettings.ResponseQueue, body: responseBytes); } private int GetRetryCount(BasicDeliverEventArgs ea) { if (ea.BasicProperties?.Headers != null && ea.BasicProperties.Headers.TryGetValue("x-retry-count", out var retryObj)) { return Convert.ToInt32(retryObj); } return 0; } }