Packages
step_flow
1.4.2-rc2
1.9.0-rc2
1.9.0-rc1
1.9.0-rc0
1.8.2
1.8.1
1.8.1-rc8
1.8.1-rc7
1.8.1-rc6
1.8.1-rc5
1.8.1-rc4
1.8.1-rc3
1.8.1-rc2
1.8.1-rc1
1.8.1-rc0
1.8.0
1.8.0-rc3
1.8.0-rc2
1.8.0-rc1
1.8.0-rc0
1.7.3
1.7.3-rc4
1.7.3-rc3
1.7.3-rc2
1.7.3-rc1
1.7.3-rc0
1.7.2
1.7.2-rc4
1.7.2-rc3
1.7.2-rc2
1.7.2-rc1
1.7.2-rc0
1.7.1
1.7.0
1.7.0-rc1
1.7.0-rc0
1.6.1
1.6.1-rc1
1.6.1-rc0
1.6.0
1.6.0-rc9
1.6.0-rc8
1.6.0-rc7
1.6.0-rc6
1.6.0-rc5
1.6.0-rc4
1.6.0-rc3
1.6.0-rc20
1.6.0-rc2
1.6.0-rc19
1.6.0-rc18
1.6.0-rc17
1.6.0-rc16
1.6.0-rc15
1.6.0-rc14
1.6.0-rc13
1.6.0-rc12
1.6.0-rc11
1.6.0-rc10
1.6.0-rc1
1.5.0
1.5.0-rc1
1.4.2-rc2
1.4.2-rc1
1.4.1
1.4.1-rc1
1.4.0
1.4.0-rc4
1.4.0-rc3
1.4.0-rc2
1.4.0-rc1
1.3.1
1.3.0
1.3.0-rc
1.2.0
1.1.0
1.0.0
1.0.0-rc9
1.0.0-rc8
1.0.0-rc7
1.0.0-rc6
1.0.0-rc5
1.0.0-rc1
0.2.13
0.2.12
0.2.11
0.2.10
0.2.9
0.2.8
0.2.7
0.2.6
0.2.5
0.2.4
0.2.3
0.2.2
0.2.1
0.2.0
0.1.8
0.1.7
0.1.6
0.1.5
0.1.4
0.1.3
0.1.2
0.1.1
0.1.0
0.0.24
0.0.23
0.0.22
0.0.21
0.0.20
0.0.19
0.0.18
0.0.17
0.0.16
0.0.15
0.0.14
0.0.13
0.0.12
0.0.11
0.0.10
0.0.9
0.0.8
0.0.7
0.0.6
0.0.4
0.0.3
0.0.2
0.0.1
Step flow manager for Elixir applications
Current section
Files
Jump to
Current section
Files
lib/step_flow/amqp/connection.ex
defmodule StepFlow.Amqp.Connection do
require Logger
@moduledoc false
@submit_exchange "job_submit"
use GenServer
alias StepFlow.Amqp.Helpers
alias StepFlow.Jobs
alias StepFlow.Metrics.JobInstrumenter
alias StepFlow.Workflows
def child_spec(_) do
%{
id: StepFlow.Amqp.Connection,
start: {StepFlow.Amqp.Connection, :start_link, []},
type: :worker
}
end
def start_link do
GenServer.start_link(__MODULE__, :ok, name: __MODULE__)
end
def consume(queue, callback) do
GenServer.cast(__MODULE__, {:consume, queue, callback})
end
def publish(queue, message, options, exchange \\ @submit_exchange) do
GenServer.cast(__MODULE__, {:publish, exchange, queue, message, options})
end
def disconnect do
GenServer.call(__MODULE__, :disconnect)
end
# def publish_json(queue, message) do
# publish(queue, message |> Jason.encode!())
# end
@impl true
def init(:ok) do
Logger.warn("#{__MODULE__} init")
rabbitmq_connect()
end
@impl true
def handle_call(:disconnect, _from, state) do
state = rabbitmq_disconnect(state)
{:reply, :ok, state}
end
@impl true
def handle_cast({:publish, exchange, queue, message, options}, state) do
Logger.info(
"#{__MODULE__}: publish message on exchange #{exchange} and queue #{queue}: #{message}"
)
AMQP.Basic.publish(state.channel, exchange, queue, message, options)
{:noreply, state}
end
@impl true
def handle_info({:DOWN, _, :process, _pid, _reason}, state) do
state =
if state.auto_restart do
{:ok, state} = rabbitmq_connect()
state
else
state
end
{:noreply, state}
end
@impl true
# Confirmation sent by the broker after registering this process as a consumer
def handle_info({:basic_consume_ok, %{consumer_tag: _consumer_tag}}, state) do
{:noreply, state}
end
@impl true
def handle_info(
{:basic_deliver, payload, %{delivery_tag: tag, redelivered: _redelivered} = _headers},
state
) do
case payload |> Jason.decode() do
{:ok, data} ->
job_id = Map.get(data, "job_id")
case job_id do
nil ->
Logger.warning("Cannot get job_id from message #{tag}: #{payload}")
job_id ->
handle_message_with_job_id(state, tag, payload, job_id)
end
{:error, _} ->
Logger.error("Payload is not a json: #{payload}")
AMQP.Basic.ack(state.channel, tag)
end
{:noreply, state}
end
@impl true
def terminate(_reason, state) do
rabbitmq_disconnect(state)
end
defp rabbitmq_connect do
url = Helpers.get_amqp_connection_url()
options = Helpers.get_amqp_connection_options()
case AMQP.Connection.open(url, options) do
{:ok, connection} ->
init_amqp_connection(connection)
{:error, message} ->
Logger.error("#{__MODULE__}: unable to connect to: #{url}, reason: #{inspect(message)}")
# Reconnection loop
:timer.sleep(10_000)
rabbitmq_connect()
end
end
defp init_amqp_connection(connection) do
Process.monitor(connection.pid)
{:ok, channel} = AMQP.Channel.open(connection)
# AMQP.Queue.declare(channel, queue)
# Logger.warn("#{__MODULE__}: connected to queue #{queue}")
AMQP.Exchange.topic(channel, @submit_exchange,
durable: true,
arguments: [{"alternate-exchange", :longstr, "job_queue_not_found"}]
)
AMQP.Exchange.fanout(channel, "job_queue_not_found", durable: true)
AMQP.Queue.declare(channel, "job_queue_not_found")
AMQP.Queue.bind(channel, "job_queue_not_found", "job_queue_not_found")
AMQP.Basic.consume(channel, "job_queue_not_found")
{:ok, %{channel: channel, connection: connection, auto_restart: true}}
end
defp handle_message_with_job_id(channel, tag, payload, job_id) do
case Jobs.get_job(job_id) do
nil ->
AMQP.Basic.reject(channel.channel, tag, requeue: true)
job ->
Logger.error("Job queue not found #{inspect(payload)}")
JobInstrumenter.inc(:step_flow_jobs_error, job.name)
description = "No worker is started with this queue name."
{:ok, job_status} = Jobs.Status.set_job_status(job_id, :error, %{message: description})
Workflows.Status.define_workflow_status(
job.workflow_id,
:queue_not_found,
job_status
)
Workflows.notification_from_job(job_id, description)
AMQP.Basic.ack(channel.channel, tag)
end
end
defp rabbitmq_disconnect(state) do
Logger.warn("#{__MODULE__}: Closing AMQP connection...")
AMQP.Connection.close(state.connection)
state
|> Map.replace(:auto_restart, false)
end
end