/
idalynnn
/
api
Обзор
Документация
Войти
/
idalynnn
/
api
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
lab4/server/Program.cs
287 строк
10 KB
idalynnn
logs
16 дек 2025, 16:36
16 дек 2025, 16:36
567b48c
Код
Авторство
О чём код?
using System.Text; using System.Text.Json; using RabbitMQ.Client; using RabbitMQ.Client.Events; const string ApiKey = "lab4-secret-key"; // <--- твоя "аутентификация" const string RequestsQueue = "api.requests"; const string RetryQueue = "api.requests.retry"; const string DlqQueue = "api.requests.dlq"; const string DlqExchange = "api.dlx"; const int MaxRetries = 3; var db = new InMemoryDb(); var idem = new IdempotencyStore(); var factory = new ConnectionFactory { HostName = "localhost", UserName = "guest", Password = "guest", DispatchConsumersAsync = true }; using var conn = factory.CreateConnection(); using var ch = conn.CreateModel(); Console.WriteLine("Connected to RabbitMQ"); // DLQ exchange + queue ch.ExchangeDeclare(DlqExchange, ExchangeType.Direct, durable: true); ch.QueueDeclare(DlqQueue, durable: true, exclusive: false, autoDelete: false); ch.QueueBind(DlqQueue, DlqExchange, routingKey: "dlq"); // Retry queue: TTL -> возвращает в requests var retryArgs = new Dictionary<string, object> { ["x-dead-letter-exchange"] = "", ["x-dead-letter-routing-key"] = RequestsQueue, ["x-message-ttl"] = 5000 // 5 секунд задержки перед повтором }; ch.QueueDeclare(RetryQueue, durable: true, exclusive: false, autoDelete: false, arguments: retryArgs); // Requests queue с DLX (фатальные ошибки -> DLQ) var reqArgs = new Dictionary<string, object> { ["x-dead-letter-exchange"] = DlqExchange, ["x-dead-letter-routing-key"] = "dlq" }; ch.QueueDeclare(RequestsQueue, durable: true, exclusive: false, autoDelete: false, arguments: reqArgs); ch.BasicQos(0, 1, false); Console.WriteLine("Server started. Waiting for messages..."); var consumer = new AsyncEventingBasicConsumer(ch); consumer.Received += async (_, ea) => { var body = ea.Body.ToArray(); var json = Encoding.UTF8.GetString(body); ApiRequest? req = null; try { req = JsonSerializer.Deserialize<ApiRequest>(json); if (req == null) throw new Exception("Invalid JSON"); Console.WriteLine($"[REQ] id={req.Id} action={req.Action} version={req.Version}"); // Идемпотентность: если id уже обработан — вернём сохранённый ответ var saved = idem.GetResponseJson(req.Id); if (saved != null) { PublishResponse(ch, ea, saved); ch.BasicAck(ea.DeliveryTag, false); return; } // Auth (API-key) if (req.Auth != ApiKey) { var resp = new ApiResponse { CorrelationId = req.Id, Status = "error", Error = "Unauthorized: invalid api-key" }; var respJson = JsonSerializer.Serialize(resp); idem.SaveResponseJson(req.Id, respJson); PublishResponse(ch, ea, respJson); ch.BasicAck(ea.DeliveryTag, false); return; } // Обработка action var response = Handle(req, db); response.CorrelationId = req.Id; Console.WriteLine($"[RES] id={req.Id} status={response.Status}"); var outJson = JsonSerializer.Serialize(response); idem.SaveResponseJson(req.Id, outJson); PublishResponse(ch, ea, outJson); ch.BasicAck(ea.DeliveryTag, false); } catch (Exception ex) { // retry / dlq var retries = GetRetryCount(ea.BasicProperties); Console.WriteLine($"Error: {ex.Message}. Retry={retries}"); if (retries < MaxRetries) { var props = ch.CreateBasicProperties(); props.Persistent = true; props.CorrelationId = ea.BasicProperties?.CorrelationId; props.ReplyTo = ea.BasicProperties?.ReplyTo; props.Headers = ea.BasicProperties?.Headers ?? new Dictionary<string, object>(); props.Headers["x-retry-count"] = retries + 1; ch.BasicPublish("", RetryQueue, props, body); ch.BasicAck(ea.DeliveryTag, false); } else { // фатально -> DLQ через reject (DLX настроен на requests) ch.BasicReject(ea.DeliveryTag, requeue: false); } } await Task.CompletedTask; }; ch.BasicConsume(RequestsQueue, autoAck: false, consumer: consumer); Console.ReadLine(); static int GetRetryCount(IBasicProperties? props) { if (props?.Headers == null) return 0; if (!props.Headers.TryGetValue("x-retry-count", out var v)) return 0; try { if (v is byte[] bytes) return int.Parse(Encoding.UTF8.GetString(bytes)); return Convert.ToInt32(v); } catch { return 0; } } static void PublishResponse(IModel ch, BasicDeliverEventArgs ea, string respJson) { var replyTo = ea.BasicProperties.ReplyTo; if (string.IsNullOrWhiteSpace(replyTo)) { // fallback: если replyTo нет — можно было бы слать в api.responses, но мы используем reply queue return; } var props = ch.CreateBasicProperties(); props.CorrelationId = ea.BasicProperties.CorrelationId; var bytes = Encoding.UTF8.GetBytes(respJson); ch.BasicPublish("", replyTo, props, bytes); } static ApiResponse Handle(ApiRequest req, InMemoryDb db) { // data приходит как JsonElement — удобно распарсить так: var dataJson = JsonSerializer.Serialize(req.Data); using var doc = JsonDocument.Parse(dataJson); var root = doc.RootElement; // v2: tasks могут иметь priority var isV2 = req.Version.Equals("v2", StringComparison.OrdinalIgnoreCase); return req.Action switch { "create_project" => CreateProject(root, db), "list_projects" => new ApiResponse { Status = "ok", Data = db.Projects.Values.ToList() }, "get_project" => GetProject(root, db), "update_project" => UpdateProject(root, db), "delete_project" => DeleteProject(root, db), "create_task" => CreateTask(root, db, isV2), "list_tasks" => new ApiResponse { Status = "ok", Data = db.Tasks.Values.ToList() }, "get_task" => GetTask(root, db), "update_task" => UpdateTask(root, db, isV2), "delete_task" => DeleteTask(root, db), _ => new ApiResponse { Status = "error", Error = $"Unknown action: {req.Action}" } }; } static ApiResponse CreateProject(JsonElement root, InMemoryDb db) { var name = root.GetProperty("name").GetString() ?? "Untitled"; var id = ++db.ProjectSeq; var p = new ProjectDto(id, name); db.Projects[id] = p; return new ApiResponse { Status = "ok", Data = p }; } static ApiResponse GetProject(JsonElement root, InMemoryDb db) { var id = root.GetProperty("id").GetInt32(); return db.Projects.TryGetValue(id, out var p) ? new ApiResponse { Status = "ok", Data = p } : new ApiResponse { Status = "error", Error = "Project not found" }; } static ApiResponse UpdateProject(JsonElement root, InMemoryDb db) { var id = root.GetProperty("id").GetInt32(); var name = root.GetProperty("name").GetString() ?? "Untitled"; if (!db.Projects.ContainsKey(id)) return new ApiResponse { Status = "error", Error = "Project not found" }; var p = new ProjectDto(id, name); db.Projects[id] = p; return new ApiResponse { Status = "ok", Data = p }; } static ApiResponse DeleteProject(JsonElement root, InMemoryDb db) { var id = root.GetProperty("id").GetInt32(); var ok = db.Projects.Remove(id); // удалим связанные задачи foreach (var tid in db.Tasks.Values.Where(t => t.ProjectId == id).Select(t => t.Id).ToList()) db.Tasks.Remove(tid); return ok ? new ApiResponse { Status = "ok", Data = new { deleted = id } } : new ApiResponse { Status = "error", Error = "Project not found" }; } static ApiResponse CreateTask(JsonElement root, InMemoryDb db, bool isV2) { var projectId = root.GetProperty("projectId").GetInt32(); if (!db.Projects.ContainsKey(projectId)) return new ApiResponse { Status = "error", Error = "Project not found for task" }; var title = root.GetProperty("title").GetString() ?? "Untitled"; var completed = root.TryGetProperty("completed", out var c) && c.GetBoolean(); int? priority = null; if (isV2 && root.TryGetProperty("priority", out var pr) && pr.ValueKind == JsonValueKind.Number) priority = pr.GetInt32(); var id = ++db.TaskSeq; var t = new TaskDto(id, projectId, title, completed, priority); db.Tasks[id] = t; return new ApiResponse { Status = "ok", Data = t }; } static ApiResponse GetTask(JsonElement root, InMemoryDb db) { var id = root.GetProperty("id").GetInt32(); return db.Tasks.TryGetValue(id, out var t) ? new ApiResponse { Status = "ok", Data = t } : new ApiResponse { Status = "error", Error = "Task not found" }; } static ApiResponse UpdateTask(JsonElement root, InMemoryDb db, bool isV2) { var id = root.GetProperty("id").GetInt32(); if (!db.Tasks.TryGetValue(id, out var old)) return new ApiResponse { Status = "error", Error = "Task not found" }; var title = root.TryGetProperty("title", out var tt) ? tt.GetString() ?? old.Title : old.Title; var completed = root.TryGetProperty("completed", out var cc) ? cc.GetBoolean() : old.Completed; int? priority = old.Priority; if (isV2 && root.TryGetProperty("priority", out var pr) && pr.ValueKind == JsonValueKind.Number) priority = pr.GetInt32(); var t = new TaskDto(id, old.ProjectId, title, completed, priority); db.Tasks[id] = t; return new ApiResponse { Status = "ok", Data = t }; } static ApiResponse DeleteTask(JsonElement root, InMemoryDb db) { var id = root.GetProperty("id").GetInt32(); var ok = db.Tasks.Remove(id); return ok ? new ApiResponse { Status = "ok", Data = new { deleted = id } } : new ApiResponse { Status = "error", Error = "Task not found" }; }