/
egopen
/
Lab1
Обзор
Документация
Войти
/
egopen
/
Lab1
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
Lab4/library/BackWorkers/UserWorker.cs
171 строка
8 KB
Egopen
lab 4
01 ноя 2025, 18:03
01 ноя 2025, 18:03
78d017c
Код
Авторство
О чём код?
using library.Controllers; using library.DTO; using library.QueueMessages; using library.Services; using library.Services.library.Services; using Microsoft.AspNetCore.Mvc; using RabbitMQ.Client; using RabbitMQ.Client.Events; using System.Security.Claims; using System.Text; using System.Text.Json; using System.Text.Json.Serialization; using System.Threading.Channels; namespace library.BackWorkers { public class UserWorker : IHostedService { private readonly TokenService _tokenService; readonly IConnection _conn; IServiceProvider _serviceProvider; IChannel _channel; delegate Task<byte[]> Handler(byte[] body, IReadOnlyBasicProperties props); Dictionary<string, Handler> _handlers; private readonly ILogger<UserWorker> _log; public UserWorker(IConnection connection, IServiceProvider serviceProvider, TokenService tokenService, ILogger<UserWorker> log) { _conn = connection; _serviceProvider = serviceProvider; _tokenService = tokenService; _log = log; _handlers = new Dictionary<string, Handler> { ["get_by_id.request"] = ProcessGetUserById, ["get_all.request"] = ProcessGetAllUsers, ["delete.request"] = ProcessDeleteUser }; } public async Task StartAsync(CancellationToken token) { _channel = await _conn.CreateChannelAsync(); await _channel.ExchangeDeclareAsync("Users", ExchangeType.Direct); await _channel.QueueDeclareAsync("users.get_by_id", durable: true, exclusive: false, autoDelete: false); await _channel.QueueDeclareAsync("users.get_all", durable: true, exclusive: false, autoDelete: false); await _channel.QueueDeclareAsync("users.delete", durable: true, exclusive: false, autoDelete: false); await _channel.QueueBindAsync("users.get_by_id", "Users", "get_by_id.request"); await _channel.QueueBindAsync("users.get_all", "Users", "get_all.request"); await _channel.QueueBindAsync("users.delete", "Users", "delete.request"); await ExecuteAsync(token); } public async Task ExecuteAsync(CancellationToken token) { try { var consumer = new AsyncEventingBasicConsumer(_channel); consumer.ReceivedAsync += async (ch, ea) => { 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 = "text/plain"; props.DeliveryMode = DeliveryModes.Persistent; props.Expiration = "36000000"; await _channel.BasicPublishAsync("Users", ea.RoutingKey.Replace(".request", ".response"), true, props, res); }; foreach (var queue in _handlers.Keys) { await _channel.BasicConsumeAsync($"users.{queue.Replace(".request", "")}", autoAck: false, consumer); } } catch (Exception ex) { _log.LogError(ex.ToString()); } } public async Task StopAsync(CancellationToken token) { await _channel.CloseAsync(); await _conn.CloseAsync(); await _channel.DisposeAsync(); await _conn.DisposeAsync(); } private async Task<byte[]> ProcessGetUserById(byte[] body, IReadOnlyBasicProperties props) { using var scope = _serviceProvider.CreateScope(); var userService = scope.ServiceProvider.GetRequiredService<UserService>(); string jsonString = Encoding.UTF8.GetString(body); var req = JsonSerializer.Deserialize<BaseMessage>(jsonString); if (req is null || req.AuthKey is null) { return ErrorResponse("Error no auth"); } var payload = req.Data as GetUserByIdPayload; if (payload is null || payload.UserId is null) { return ErrorResponse("Error no data"); } var claims = _tokenService.GetClaimsFromExpiredToken(req.AuthKey); var role = claims.FirstOrDefault(c => c.Type == "role"); if (role is null || role.Value != "Admin") { return ErrorResponse("Not authorized"); } var user = await userService.GetUserById(payload.UserId.Value); return Encoding.UTF8.GetBytes(jsonString); } private async Task<byte[]> ProcessGetAllUsers(byte[] body, IReadOnlyBasicProperties props) { using var scope = _serviceProvider.CreateScope(); var userService = scope.ServiceProvider.GetRequiredService<UserService>(); string jsonString = Encoding.UTF8.GetString(body); var req = JsonSerializer.Deserialize<BaseMessage>(jsonString); if (req is null || req.AuthKey is null) { return Encoding.UTF8.GetBytes("Error"); } var payload = req.Data as GetAllUsersPayload; if (payload?.Page < 1 || payload?.PageSize < 1 || payload?.PageSize > 100) { return ErrorResponse("Invalid page parameters"); } var claims = _tokenService.GetClaimsFromExpiredToken(req.AuthKey); var role = claims.FirstOrDefault(c => c.Type == "role"); if (role is null || role.Value != "Admin") { return Encoding.UTF8.GetBytes("Not authorized"); } var user = await userService.GetAllUsers(payload is null ? 1: payload.Page, payload is null ? 20 : payload.PageSize); return Encoding.UTF8.GetBytes(JsonSerializer.Serialize(user)); } private async Task<byte[]> ProcessDeleteUser(byte[] body, IReadOnlyBasicProperties props) { using var scope = _serviceProvider.CreateScope(); var userService = scope.ServiceProvider.GetRequiredService<UserService>(); string jsonString = Encoding.UTF8.GetString(body); var req = JsonSerializer.Deserialize<BaseMessage>(jsonString); if (req is null || req.AuthKey is null) { return ErrorResponse("Error no auth"); } var payload = req.Data as GetUserByIdPayload; if (payload is null || payload.UserId is null) { return ErrorResponse("Error no data"); } var claims = _tokenService.GetClaimsFromExpiredToken(req.AuthKey); var role = claims.FirstOrDefault(c => c.Type == "role"); if (role is null || role.Value != "Admin") { return ErrorResponse("Not authorized"); } var user = await userService.DeleteUser(payload.UserId.Value); return Encoding.UTF8.GetBytes(jsonString); } private class GetUserByIdPayload { public Guid? UserId { get; set; } } private class GetAllUsersPayload { public int Page { get; set; } public int PageSize { get; set; } } private class DeleteUserPayload { public Guid? UserId { get; set; } } private byte[] ErrorResponse(string message) { var error = new { error = message }; return Encoding.UTF8.GetBytes(JsonSerializer.Serialize(error)); } } }