Current section
Files
Jump to
Current section
Files
lib/mammoth.ex
defmodule Mammoth do
@moduledoc ~S"""
Mammoth: A STOMP client.
Use `Mammoth` as the primary API and `Mammoth.Message` for working with received messages.
## Example
{:ok, pid} = Mammoth.start_link
Mammoth.connect(pid, {127,0,0,1}, 61613, "admin", "admin")
callback = fn m -> IO.puts(inspect(m)) end
Mammoth.subscribe(pid, "foo.bar", callback)
Mammoth.disconnect(pid)
For more control, use pattern matching in the callback:
callback = fn
%Mammoth.Message{command: :message, headers: headers, body: body} ->
Logger.info(["Received MESSAGE", "\nheaders: ", inspect(headers), "\nbody: ", inspect(body)])
%Mammoth.Message{command: :error, headers: headers, body: body} ->
Logger.error(["Received ERROR", "\nheaders: ", inspect(headers), "\nbody: ", inspect(body)])
%Mammoth.Message{command: cmd, headers: headers, body: body} ->
Logger.error([
"Received unknown command: ", cmd,
"\nheaders: ",
inspect(headers),
"\nbody: ",
inspect(body)
])
end
## Starting in a supervision tree
children = [
worker(Mammoth, [%{}, [name: Mammoth]])
]
"""
use GenServer
require Logger
alias Mammoth.{Message, Receiver, Socket, Subscriber}
def init(args) do
{
:ok,
args
|> Map.put_new(:socket, nil)
|> Map.put_new(:receiver, nil)
|> Map.put_new(:subscriber, nil)
}
end
def start_link(state \\ %{}, opts \\ []) do
GenServer.start_link(__MODULE__, state, opts)
end
@doc """
Connect to server.
`host` must be `inet:socket_address() | inet:hostname()`, for example `{127,0,0,1}`.
"""
def connect(pid, host, port, login, password) do
GenServer.call(pid, {:connect, host, port, login, password})
end
@doc """
Subscribe to a queue and register a callback for received messages.
"""
def subscribe(pid, destination, callback) do
GenServer.call(pid, {:subscribe, destination, callback})
end
@doc """
Unsubscribe from a queue.
"""
def unsubscribe(pid, destination) do
GenServer.call(pid, {:unsubscribe, destination})
end
@doc """
Receive messages from the TCP socket.
Is called automatically when necessary. Should not be called manually.
"""
def receive(pid, message) do
GenServer.call(pid, {:receive, message})
end
@doc """
Disconnect from server.
"""
def disconnect(pid) do
GenServer.call(pid, :disconnect)
end
def handle_call(
{:receive, message = %Message{command: :message}},
_from,
state = %{subscriber: subscriber}
) do
{:ok, destination} = Message.get_header(message, "destination")
%{callback: callback} = Subscriber.get_subscription(subscriber, destination)
callback.(message)
{:reply, :ok, state}
end
def handle_call({:connect, host, port, login, password}, _from, state) do
# todo: handle {:error, :econnrefused} response
{:ok, socket} = Socket.connect(host, port)
Socket.send(socket, connect_message(login, password))
{:ok, response} = Socket.receive(socket)
{:ok, subscriber} = Subscriber.start_link()
{:ok, receiver} = Receiver.start_link(%{socket: socket, consumer: self()})
Receiver.listen(receiver)
{:ok, message, _} = Message.parse(response)
case message do
%Message{command: :connected} ->
{:reply, socket, %{state | socket: socket, subscriber: subscriber, receiver: receiver}}
_ ->
Logger.warn(
"[Mammoth] failed to connect to host: #{inspect(host)}, port: #{port}, login: #{login}, reason: #{inspect(message)}"
)
{:reply, {:error, message}, state}
end
end
@doc """
Disconnects from server.
Returns `{:ok, :disconnected}` or `{:error, :disconnect_failed, message}`.
"""
def handle_call(:disconnect, _from, %{
socket: socket,
subscriber: subscriber,
receiver: receiver
}) do
receipt_id = Enum.random(1000..1_000_000)
Socket.send(socket, disconnect_message(receipt_id))
Receiver.stop(receiver)
Subscriber.stop(subscriber)
case Socket.receive(socket) do
{:error, :closed} ->
{:reply, {:ok, :disconnected}, %{}}
{:error, reason} ->
{:reply, {:error, reason}, %{}}
{:ok, response} ->
Logger.debug(response)
{:ok, response_message, _} = Message.parse(response)
if response_message.command == :receipt &&
Message.has_header(response_message, {"receipt-id", receipt_id}) do
{:reply, {:ok, :disconnected}, %{}}
else
{:reply, {:error, :disconnect_failed, response_message}, %{}}
end
end
end
def handle_call(
{:subscribe, destination, callback},
_from,
state = %{socket: socket, subscriber: subscriber}
) do
{:ok, entry} = Subscriber.subscribe(subscriber, destination, callback)
message = subscribe_message(destination, entry.id)
Socket.send(socket, message)
{:reply, {:ok, entry}, state}
end
def handle_call(
{:unsubscribe, destination},
_from,
state = %{socket: socket, subscriber: subscriber}
) do
{:ok, entry} = Subscriber.unsubscribe(subscriber, destination)
message = unsubscribe_message(entry.id)
Socket.send(socket, message)
{:reply, :ok, state}
end
defp connect_message(login, password) do
%Message{
command: :connect,
headers: [
{"accept-version", "1.2"},
{"host", "localhost"},
{"login", login},
{"passcode", password}
]
}
end
defp disconnect_message(receipt_id) do
%Message{
command: :disconnect,
headers: [
{"receipt-id", receipt_id}
]
}
end
defp subscribe_message(destination, id) do
%Message{
command: :subscribe,
headers: [
{"destination", destination},
{"ack", "auto"},
{"id", id}
]
}
end
defp unsubscribe_message(id) do
%Message{
command: :unsubscribe,
headers: [
{"id", id}
]
}
end
end