/
ArtSerg
/
IntegrationProject
Обзор
Документация
Войти
/
ArtSerg
/
IntegrationProject
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
lab4
Integration1/Services/MessageHandlerService.cs
658 строк
24 KB
artS3rg
Добавлена интеграция RabbitMQ для асинхронного обмена сообщениями
14 дек 2025, 21:50
14 дек 2025, 21:50
a02de12
Код
Авторство
О чём код?
using Integration1.Data; using Integration1.Dtos; using Integration1.Messaging; using Integration1.Models; using Microsoft.EntityFrameworkCore; using Microsoft.IdentityModel.Tokens; using System.IdentityModel.Tokens.Jwt; using System.Security.Claims; using System.Security.Cryptography; using System.Text; using System.Text.Json; using System.Text.Json.Serialization; namespace Integration1.Services { public class MessageHandlerService { private readonly AppDbContext _db; private readonly JwtService _jwt; private readonly IConfiguration _config; private readonly ILogger<MessageHandlerService> _logger; public MessageHandlerService( AppDbContext db, JwtService jwt, IConfiguration config, ILogger<MessageHandlerService> logger) { _db = db; _jwt = jwt; _config = config; _logger = logger; } public async Task<ResponseMessage> HandleMessage(RequestMessage request) { _logger.LogInformation("Handling message: Id={MessageId}, Action={Action}, Auth present: {HasAuth}", request.Id, request.Action, !string.IsNullOrEmpty(request.Auth)); // Проверка идемпотентности var processed = await _db.ProcessedMessages.FindAsync(request.Id); if (processed != null) { var dataPreview = processed.ResponseData != null && processed.ResponseData.Length > 100 ? processed.ResponseData.Substring(0, 100) : processed.ResponseData; _logger.LogInformation("Duplicate message detected: {MessageId}, Action={Action}, returning cached response. Cached data preview: {Data}", request.Id, processed.Action, dataPreview); try { object? data = null; if (!string.IsNullOrEmpty(processed.ResponseData)) { data = JsonSerializer.Deserialize<object>(processed.ResponseData); } var cachedResponse = new ResponseMessage { CorrelationId = request.Id, // Используем ID запроса как correlation_id Status = "ok", Data = data }; _logger.LogInformation("Returning cached response for message: {MessageId}, CorrelationId: {CorrelationId}", request.Id, cachedResponse.CorrelationId); return cachedResponse; } catch (Exception ex) { _logger.LogWarning(ex, "Failed to deserialize cached response for message: {MessageId}, will process as new", request.Id); // Если не удалось десериализовать, продолжаем обработку как новое сообщение } } else { _logger.LogInformation("New message: {MessageId}, Action={Action}, will process", request.Id, request.Action); } // Проверка аутентификации if (!ValidateAuth(request.Auth)) { return new ResponseMessage { CorrelationId = request.Id, Status = "error", Error = "Unauthorized: Invalid API key" }; } ResponseMessage response; try { response = request.Action.ToLower() switch { "register" => await HandleRegister(request), "login" => await HandleLogin(request), "get_tasks" => await HandleGetTasks(request), "create_task" => await HandleCreateTask(request), "update_task" => await HandleUpdateTask(request), "delete_task" => await HandleDeleteTask(request), _ => new ResponseMessage { CorrelationId = request.Id, Status = "error", Error = $"Unknown action: {request.Action}" } }; } catch (Exception ex) { _logger.LogError(ex, "Error handling action: {Action}", request.Action); response = new ResponseMessage { CorrelationId = request.Id, Status = "error", Error = ex.Message }; } // Сохраняем обработанное сообщение для идемпотентности if (response.Status == "ok") { try { // Проверяем, не существует ли уже запись (на случай race condition) var existing = await _db.ProcessedMessages.FindAsync(request.Id); if (existing == null) { _db.ProcessedMessages.Add(new ProcessedMessage { MessageId = request.Id, Action = request.Action, ResponseData = JsonSerializer.Serialize(response.Data), ProcessedAt = DateTime.UtcNow }); await _db.SaveChangesAsync(); } } catch (Exception ex) { // Игнорируем ошибки при сохранении идемпотентности (например, если запись уже существует) _logger.LogWarning(ex, "Failed to save processed message for idempotency: {MessageId}", request.Id); } } return response; } private bool ValidateAuth(string? auth) { if (string.IsNullOrEmpty(auth)) return false; var apiKey = _config["Internal:ApiKey"]; return auth == apiKey; } private async Task<ResponseMessage> HandleRegister(RequestMessage request) { RegisterDto? data = null; try { string jsonString; if (request.Data is JsonElement jsonElement) { jsonString = jsonElement.GetRawText(); _logger.LogInformation("Register data is JsonElement: {Json}", jsonString); } else if (request.Data is Dictionary<string, object> dict) { jsonString = JsonSerializer.Serialize(dict); _logger.LogInformation("Register data is Dictionary: {Json}", jsonString); } else { jsonString = JsonSerializer.Serialize(request.Data); _logger.LogInformation("Register data is other type: {Type}, {Json}", request.Data?.GetType().Name, jsonString); } var options = new JsonSerializerOptions { PropertyNameCaseInsensitive = true }; data = JsonSerializer.Deserialize<RegisterDto>(jsonString, options); _logger.LogInformation("Deserialized RegisterDto: Username={Username}, Password={Password}, JSON was: {Json}", data?.Username, data?.Password != null ? "***" : null, jsonString); } catch (Exception ex) { _logger.LogWarning(ex, "Failed to deserialize register data. Data type: {Type}", request.Data?.GetType().Name); } if (data == null || string.IsNullOrEmpty(data.Username) || string.IsNullOrEmpty(data.Password)) { _logger.LogWarning("Invalid registration data. Data is null: {IsNull}, Username: {Username}, Password: {Password}", data == null, data?.Username, data?.Password != null ? "present" : "missing"); return new ResponseMessage { CorrelationId = request.Id, Status = "error", Error = "Invalid registration data" }; } if (await _db.Users.AnyAsync(u => u.Username == data.Username)) { return new ResponseMessage { CorrelationId = request.Id, Status = "error", Error = "Username already exists" }; } var user = new User { Username = data.Username, PasswordHash = HashPassword(data.Password) }; _db.Users.Add(user); await _db.SaveChangesAsync(); return new ResponseMessage { CorrelationId = request.Id, Status = "ok", Data = new { message = "User registered successfully", userId = user.Id } }; } private async Task<ResponseMessage> HandleLogin(RequestMessage request) { LoginDto? data = null; try { string jsonString; if (request.Data is JsonElement jsonElement) { jsonString = jsonElement.GetRawText(); _logger.LogInformation("Login data is JsonElement: {Json}", jsonString); } else { jsonString = JsonSerializer.Serialize(request.Data); _logger.LogInformation("Login data is other type: {Type}, {Json}", request.Data?.GetType().Name, jsonString); } var options = new JsonSerializerOptions { PropertyNameCaseInsensitive = true }; data = JsonSerializer.Deserialize<LoginDto>(jsonString, options); _logger.LogInformation("Deserialized LoginDto: Username={Username}, Password={Password}", data?.Username, data?.Password != null ? "***" : null); } catch (Exception ex) { _logger.LogWarning(ex, "Failed to deserialize login data"); } if (data == null || string.IsNullOrEmpty(data.Username) || string.IsNullOrEmpty(data.Password)) { return new ResponseMessage { CorrelationId = request.Id, Status = "error", Error = "Invalid login data" }; } var user = await _db.Users.FirstOrDefaultAsync(u => u.Username == data.Username); if (user == null || user.PasswordHash != HashPassword(data.Password)) { return new ResponseMessage { CorrelationId = request.Id, Status = "error", Error = "Invalid username or password" }; } var token = _jwt.GenerateToken(user); return new ResponseMessage { CorrelationId = request.Id, Status = "ok", Data = new { token } }; } private async Task<ResponseMessage> HandleGetTasks(RequestMessage request) { var userId = ExtractUserId(request); if (string.IsNullOrEmpty(userId)) { return new ResponseMessage { CorrelationId = request.Id, Status = "error", Error = "User not authenticated" }; } int page = 1; int pageSize = 10; if (request.Data != null) { Dictionary<string, object>? dataDict = null; if (request.Data is JsonElement jsonElement) { dataDict = JsonSerializer.Deserialize<Dictionary<string, object>>(jsonElement.GetRawText()); } else if (request.Data is Dictionary<string, object> dict) { dataDict = dict; } if (dataDict != null) { if (dataDict.ContainsKey("page")) int.TryParse(dataDict["page"]?.ToString(), out page); if (dataDict.ContainsKey("pageSize")) int.TryParse(dataDict["pageSize"]?.ToString(), out pageSize); } } if (page < 1) page = 1; if (pageSize < 1) pageSize = 10; if (pageSize > 100) pageSize = 100; var q = _db.Tasks.Where(t => t.UserId == userId).AsNoTracking(); var total = await q.CountAsync(); var totalPages = (int)Math.Ceiling(total / (double)pageSize); var items = await q .OrderBy(t => t.Id) .Skip((page - 1) * pageSize) .Take(pageSize) .ToListAsync(); return new ResponseMessage { CorrelationId = request.Id, Status = "ok", Data = new { items = items.Select(t => new { t.Id, t.Title, t.Description, Priority = t.Priority.ToString(), t.IsCompleted, t.UserId, t.DueDate, t.Category }), totalCount = total, page, pageSize, totalPages } }; } private async Task<ResponseMessage> HandleCreateTask(RequestMessage request) { var userId = ExtractUserId(request); if (string.IsNullOrEmpty(userId)) { return new ResponseMessage { CorrelationId = request.Id, Status = "error", Error = "User not authenticated" }; } TaskCreateDto? data = null; try { string jsonString; if (request.Data is JsonElement jsonElement) { jsonString = jsonElement.GetRawText(); _logger.LogInformation("CreateTask data is JsonElement: {Json}", jsonString); } else { jsonString = JsonSerializer.Serialize(request.Data); _logger.LogInformation("CreateTask data is other type: {Type}, {Json}", request.Data?.GetType().Name, jsonString); } var options = new JsonSerializerOptions { PropertyNameCaseInsensitive = true, Converters = { new JsonStringEnumConverter() } }; data = JsonSerializer.Deserialize<TaskCreateDto>(jsonString, options); _logger.LogInformation("Deserialized TaskCreateDto: Title={Title}", data?.Title); } catch (Exception ex) { _logger.LogWarning(ex, "Failed to deserialize task create data. Data type: {Type}", request.Data?.GetType().Name); } if (data == null || string.IsNullOrEmpty(data.Title)) { return new ResponseMessage { CorrelationId = request.Id, Status = "error", Error = "Invalid task data" }; } var task = new TaskItem { Title = data.Title, Description = data.Description, Priority = data.Priority, UserId = userId, DueDate = data.DueDate, Category = data.Category }; _db.Tasks.Add(task); await _db.SaveChangesAsync(); return new ResponseMessage { CorrelationId = request.Id, Status = "ok", Data = new { task.Id, task.Title, task.Description, Priority = task.Priority.ToString(), task.IsCompleted, task.UserId, task.DueDate, task.Category } }; } private async Task<ResponseMessage> HandleUpdateTask(RequestMessage request) { var userId = ExtractUserId(request); if (string.IsNullOrEmpty(userId)) { return new ResponseMessage { CorrelationId = request.Id, Status = "error", Error = "User not authenticated" }; } var dataDict = GetDataDictionary(request); if (dataDict == null || !dataDict.ContainsKey("id")) { return new ResponseMessage { CorrelationId = request.Id, Status = "error", Error = "Task id is required" }; } var taskId = int.Parse(dataDict["id"]!.ToString()!); var task = await _db.Tasks.FindAsync(taskId); if (task == null || task.UserId != userId) { return new ResponseMessage { CorrelationId = request.Id, Status = "error", Error = "Task not found" }; } TaskUpdateDto? data = null; try { if (request.Data is JsonElement jsonElement) { data = JsonSerializer.Deserialize<TaskUpdateDto>(jsonElement.GetRawText()); } else { var jsonString = JsonSerializer.Serialize(request.Data); data = JsonSerializer.Deserialize<TaskUpdateDto>(jsonString); } } catch (Exception ex) { _logger.LogWarning(ex, "Failed to deserialize task update data"); } if (data != null) { if (data.Title != null) task.Title = data.Title; if (data.Description != null) task.Description = data.Description; if (data.Priority.HasValue) task.Priority = data.Priority.Value; if (data.IsCompleted.HasValue) task.IsCompleted = data.IsCompleted.Value; if (data.DueDate.HasValue) task.DueDate = data.DueDate; if (data.Category != null) task.Category = data.Category; } await _db.SaveChangesAsync(); return new ResponseMessage { CorrelationId = request.Id, Status = "ok", Data = new { task.Id, task.Title, task.Description, Priority = task.Priority.ToString(), task.IsCompleted, task.UserId, task.DueDate, task.Category } }; } private async Task<ResponseMessage> HandleDeleteTask(RequestMessage request) { var userId = ExtractUserId(request); if (string.IsNullOrEmpty(userId)) { return new ResponseMessage { CorrelationId = request.Id, Status = "error", Error = "User not authenticated" }; } var dataDict = GetDataDictionary(request); if (dataDict == null || !dataDict.ContainsKey("id")) { return new ResponseMessage { CorrelationId = request.Id, Status = "error", Error = "Task id is required" }; } var taskId = int.Parse(dataDict["id"]!.ToString()!); var task = await _db.Tasks.FindAsync(taskId); if (task == null || task.UserId != userId) { return new ResponseMessage { CorrelationId = request.Id, Status = "error", Error = "Task not found" }; } _db.Tasks.Remove(task); await _db.SaveChangesAsync(); return new ResponseMessage { CorrelationId = request.Id, Status = "ok", Data = new { message = "Task deleted successfully" } }; } private Dictionary<string, object>? GetDataDictionary(RequestMessage request) { if (request.Data == null) return null; if (request.Data is JsonElement jsonElement) { return JsonSerializer.Deserialize<Dictionary<string, object>>(jsonElement.GetRawText()); } else if (request.Data is Dictionary<string, object> dict) { return dict; } else { var jsonString = JsonSerializer.Serialize(request.Data); return JsonSerializer.Deserialize<Dictionary<string, object>>(jsonString); } } private string? ExtractUserId(RequestMessage request) { // Извлекаем userId из токена в data var dataDict = GetDataDictionary(request); if (dataDict?.ContainsKey("token") == true) { var token = dataDict["token"]?.ToString(); if (!string.IsNullOrEmpty(token)) { try { var handler = new JwtSecurityTokenHandler(); var key = new SymmetricSecurityKey(Encoding.UTF8.GetBytes(_config["Jwt:Key"]!)); var validationParams = new TokenValidationParameters { ValidateIssuer = false, ValidateAudience = false, ValidateLifetime = true, ValidateIssuerSigningKey = true, IssuerSigningKey = key }; var principal = handler.ValidateToken(token, validationParams, out _); var userId = principal.FindFirst("UserId")?.Value; if (!string.IsNullOrEmpty(userId)) return userId; } catch (Exception ex) { _logger.LogWarning(ex, "Failed to validate token"); } } } // Fallback: можно передавать userId напрямую в data if (dataDict?.ContainsKey("userId") == true) { return dataDict["userId"]?.ToString(); } return null; } private static string HashPassword(string password) { using var sha = SHA256.Create(); var bytes = sha.ComputeHash(Encoding.UTF8.GetBytes(password)); return Convert.ToBase64String(bytes); } } }