/
egopen
/
Lab1
Обзор
Документация
Войти
/
egopen
/
Lab1
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
Lab4/library/BackWorkers/AuthWorker.cs
159 строк
6 KB
Egopen
lab 4
01 ноя 2025, 18:03
01 ноя 2025, 18:03
78d017c
Код
Авторство
О чём код?
using library.QueueMessages; using library.Services; using Microsoft.AspNetCore.Mvc; using RabbitMQ.Client; using RabbitMQ.Client.Events; using System.Text; using System.Text.Json; namespace library.BackWorkers { public class AuthWorker : IHostedService { readonly IConnection _conn; IServiceProvider _serviceProvider; IChannel _channel; delegate Task<byte[]> Handler(byte[] body, IReadOnlyBasicProperties props); Dictionary<string, Handler> _handlers; private readonly ILogger<AuthWorker> _log; public AuthWorker(IConnection connection, IServiceProvider serviceProvider, ILogger<AuthWorker> log) { _conn = connection; _serviceProvider = serviceProvider; _log = log; _handlers = new Dictionary<string, Handler> { ["register.request"] = Register, ["login.request"] = Login }; } public async Task StartAsync(CancellationToken token) { _channel = await _conn.CreateChannelAsync(); var dlxName = "Books.DLX"; var dlqName = "books.get_book.dlq"; await _channel.ExchangeDeclareAsync(dlxName, ExchangeType.Fanout); await _channel.QueueDeclareAsync(dlqName, durable: true, exclusive: false, autoDelete: false); await _channel.QueueBindAsync(dlqName, dlxName, ""); var queueArgs = new Dictionary<string, object> { { "x-dead-letter-exchange", dlxName }, { "x-dead-letter-routing-key", dlqName } }; await _channel.ExchangeDeclareAsync("Auth", ExchangeType.Direct); await _channel.QueueDeclareAsync("auth.register", durable: true, exclusive: false, autoDelete: false); await _channel.QueueDeclareAsync("auth.login", durable: true, exclusive: false, autoDelete: false); await _channel.QueueBindAsync("auth.register", "Auth", "register.request"); await _channel.QueueBindAsync("auth.login", "Auth", "login.request"); await ExecuteAsync(token); } public async Task ExecuteAsync(CancellationToken token) { try { var consumer = new AsyncEventingBasicConsumer(_channel); consumer.ReceivedAsync += async (ch, ea) => { try { var body = ea.Body.ToArray(); var res = await _handlers[ea.RoutingKey](body, ea.BasicProperties); await _channel.BasicAckAsync(ea.DeliveryTag, false); var props = new BasicProperties(); props.ContentType = "application/json"; props.DeliveryMode = DeliveryModes.Persistent; await _channel.BasicPublishAsync("Auth", ea.RoutingKey.Replace(".request", ".response"), true, props, res); } catch (Exception ex) { _log.LogError(ex, "Failed to process message, sending to DLQ"); await _channel.BasicNackAsync(ea.DeliveryTag, false, false); } }; foreach (var queue in _handlers.Keys) { await _channel.BasicConsumeAsync($"auth.{queue.Replace(".request", "")}", autoAck: false, consumer); } } catch (Exception ex) { _log.LogError(ex.ToString()); } } public async Task<byte[]> Register(byte[] body, IReadOnlyBasicProperties props) { using var scope = _serviceProvider.CreateScope(); var authService = scope.ServiceProvider.GetRequiredService<AuthService>(); var jsonString = Encoding.UTF8.GetString(body); var req = JsonSerializer.Deserialize<BaseMessage>(jsonString); if (req?.Data == null) return ErrorResponse("Invalid request"); var payload = JsonSerializer.Deserialize<RegisterPayload>(req.Data.ToString()); if (payload?.Name == null || payload?.Email == null || payload?.Password == null) return ErrorResponse("Invalid payload"); var token = await authService.RegisterUser(payload.Name, payload.Email, payload.Password, payload.Role); if (token == null) return ErrorResponse("User with this email already exists", 409); return Encoding.UTF8.GetBytes(JsonSerializer.Serialize(new { access_token = token })); } public async Task<byte[]> Login(byte[] body, IReadOnlyBasicProperties props) { using var scope = _serviceProvider.CreateScope(); var authService = scope.ServiceProvider.GetRequiredService<AuthService>(); var jsonString = Encoding.UTF8.GetString(body); var req = JsonSerializer.Deserialize<BaseMessage>(jsonString); if (req?.Data == null) return ErrorResponse("Invalid request"); var payload = JsonSerializer.Deserialize<LoginPayload>(req.Data.ToString()); if (payload?.Email == null || payload?.Password == null) return ErrorResponse("Invalid payload"); var token = await authService.AuthUser(payload.Email, payload.Password); return Encoding.UTF8.GetBytes(JsonSerializer.Serialize(new { access_token = token })); } public async Task StopAsync(CancellationToken token) { await _channel.CloseAsync(); await _channel.DisposeAsync(); } private byte[] ErrorResponse(string message, int statusCode = 400) { var error = new { error = message, statusCode }; return Encoding.UTF8.GetBytes(JsonSerializer.Serialize(error)); } private class RegisterPayload { public string Name { get; set; } public string Email { get; set; } public string Password { get; set; } public string Role { get; set; } } private class LoginPayload { public string Email { get; set; } public string Password { get; set; } } } }