/
count_rochester
/
rust_microservices
Обзор
Документация
Войти
/
count_rochester
/
rust_microservices
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
services/orders/src/controllers/rmq_controller.rs
77 строк
3 KB
Гусев Андрей
Added create orders logic
09 фев 2025, 20:28
09 фев 2025, 20:28
043cee5
Код
Авторство
О чём код?
use broker::{deserialize_message, get_routing_key, serialize_message, PublishType, RmqMessage}; use chrono::Utc; use common::outbox::{CreateOutboxMessagePayload, OutboxReceiverService}; use constants::{broker::ORDER_REQUEST_RK_PREFIX, exchange::order::ORDER_RK_PREFIX}; use models::exchange::{ ExchangeRequestDto, ExchangeResponseDto, ExchangeResponseStatus, OrderCartRequestPayload, }; use uuid::Uuid; use crate::{config, repository::get_store, services::OrdersService}; use super::{Error, Result}; pub struct RmqController; impl RmqController { pub async fn process_message<'a>(message: RmqMessage<'a>) -> Result<()> { let order_request_key = get_routing_key(ORDER_REQUEST_RK_PREFIX.to_string(), PublishType::Request); match message.routing_key { key if key == order_request_key => match message.data { None => Err(Error::InvalidRmqPayload)?, Some(data) => { let payload = deserialize_message::<ExchangeRequestDto<OrderCartRequestPayload>>(data) .map_err(|_| Error::InvalidRmqPayload)?; if message.retry_count.unwrap_or(0) > config().MAX_ORDER_REQUEST_RETRY_COUNT as i64 { Self::inform_error(payload.id, ORDER_RK_PREFIX.to_string()).await?; } else { Self::on_create_order(payload).await?; } } }, _ => Err(Error::RKError(message.routing_key))?, } Ok(()) } async fn on_create_order(payload: ExchangeRequestDto<OrderCartRequestPayload>) -> Result<()> { let service = OrdersService::new(); service .on_create_order(payload) .await .map_err(|e| Error::CreateOrderError(e.to_string()))?; Ok(()) } async fn inform_error(request_id: Uuid, prefix: String) -> Result<()> { let message: ExchangeResponseDto = ExchangeResponseDto { id: request_id.clone(), status: ExchangeResponseStatus::Failure, error: Some("Max attempts reached".to_string()), }; let outbox = OutboxReceiverService::new(get_store().await.clone()); let outbox_message = CreateOutboxMessagePayload { object_id: request_id.to_string(), object_modify_date: Utc::now(), routing_key: get_routing_key(prefix, PublishType::Response), content: serialize_message(message).map_err(|_| Error::SerializeMessageError)?, }; outbox .accept_new_message(outbox_message, None) .await .map_err(|e| Error::PublishOutboxError(e.to_string()))?; Ok(()) } }