Packages
step_flow
1.0.0
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 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 publish_json(queue, message) do
# publish(queue, message |> Jason.encode!())
# end
def init(:ok) do
Logger.warn("#{__MODULE__} init")
rabbitmq_connect()
end
def handle_cast({:publish, exchange, queue, message, options}, conn) do
Logger.info(
"#{__MODULE__}: publish message on exchange #{exchange} and queue #{queue}: #{message}"
)
AMQP.Basic.publish(conn.channel, exchange, queue, message, options)
{:noreply, conn}
end
def handle_info({:DOWN, _, :process, _pid, _reason}, _) do
{:ok, chan} = rabbitmq_connect()
{:noreply, chan}
end
# Confirmation sent by the broker after registering this process as a consumer
def handle_info({:basic_consume_ok, %{consumer_tag: _consumer_tag}}, chan) do
{:noreply, chan}
end
def handle_info(
{:basic_deliver, payload, %{delivery_tag: tag, redelivered: _redelivered} = _headers},
channel
) do
data =
payload
|> Jason.decode!()
job_id = Map.get(data, "job_id")
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
{:noreply, channel}
end
def terminate(_reason, state) do
AMQP.Connection.close(state.connection)
end
defp rabbitmq_connect do
url = Helpers.get_amqp_connection_url()
case AMQP.Connection.open(url) 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}}
end
end