defmodule VtexWs.Workers.NotificationWorker do @moduledoc "Worker responsável por consolidar os produtos e enviá-los para a fila SQS." alias VtexWs.{Apps, Products, Categories} alias VtexWs.Products.AppManagerProduct, as: Product import Task, only: [async: 1, await: 1] defexception [:message] @queue "product_consolidation" @doc "Retorna o nome da fila do Redis." def queue, do: @queue @doc "Retorna a url da fila na AWS." def fifo, do: System.get_env("SQS_FIFO_URL") @doc "Método que consolida os produtos e os envia para a fila." def perform( app_slug, product_sku, product_price, product_stock_in_store, retry ) when retry < 10 do IO.puts("Start processing #{product_sku}") app = Apps.get_app_by!(slug: app_slug) price = async(fn -> product_price || Products.get_price_by_app!(product_sku, app) end) stock_in_store = async(fn -> product_stock_in_store || Products.get_stock_in_store_by_app!(product_sku, app) end) product = Products.get_by_app!(product_sku, app) :ok = Categories.consolidate_all_by_app!(product.product_category_ids, app) product = %{product | app_id: app.app_id} product = %{product | price: await(price), stock_in_store: await(stock_in_store)} product = %{product | is_trash: to_string(product.price) in ["", "0"]} if Product.valid?(product) do %{status_code: 200} = send_to_sqs(product, app) IO.puts("Finished processing #{product_sku}") else args = [ app_slug, product_sku, product.price, product.stock_in_store, retry + 1 ] {:ok, _jid} = Exq.enqueue(Exq, @queue, __MODULE__, args) end end @doc "Caso as tentativas sejam maiores que 10, esse método dispara um Bugsnag." def perform( app_slug, product_sku, product_price, product_stock_in_store, retry ) do Bugsnag.report( %__MODULE__{message: "Failed to consolidate product."}, severity: :error, context: "worker#perform", metadata: %{ arguments: %{ app_slug: app_slug, product_sku: product_sku, product_price: product_price, product_stock_in_store: product_stock_in_store, retry: retry } } ) end @doc "Método que envia o produto consolidado para o SQS." def send_to_sqs(%Product{} = product, app) do fifo() |> ExAws.SQS.send_message( product |> Product.compact! |> Poison.encode!, [{:message_group_id, app.id}, {:message_deduplication_id, product.sku}] ) |> ExAws.request! end end