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, callback_pid} = Mammoth.DefaultCallbackHandler.start_link
{:ok, pid} = Mammoth.start_link(callback_pid)
Mammoth.connect(pid, {127,0,0,1}, 61613, "admin", "admin")
Mammoth.subscribe(pid, "foo.bar")
Mammoth.disconnect(pid)
## 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(callback_handler, state \\ %{}, opts \\ []) do
GenServer.start_link(
__MODULE__,
state |> Map.put_new(:callback_handler, callback_handler),
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.
"""
def subscribe(pid, destination) do
GenServer.call(pid, {:subscribe, destination})
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},
_from,
state = %{callback_handler: callback_handler}
) do
Kernel.send(callback_handler, {:mammoth, :receive_frame, 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, subscriber} = Subscriber.start_link()
{:ok, receiver} = Receiver.start_link(%{socket: socket, consumer: self()})
Receiver.listen(receiver)
{:reply, socket, %{state | socket: socket, subscriber: subscriber, receiver: receiver}}
end
@doc """
Requests disconnection from the remote server
"""
def handle_call(
:disconnect,
_from,
state = %{
socket: socket
}
) do
receipt_id = Enum.random(1000..1_000_000)
Socket.send(socket, disconnect_message(receipt_id))
{:noreply, Map.put(state, :disconnect_id, receipt_id)}
end
def handle_call(
{:subscribe, destination},
_from,
state = %{socket: socket, subscriber: subscriber}
) do
{:ok, entry} = Subscriber.subscribe(subscriber, destination)
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
def handle_cast(
:disconnected,
state = %{
subscriber: subscriber,
receiver: receiver,
callback_handler: callback_handler,
disconnect_id: _disconnect_id
}
) do
Subscriber.stop(subscriber)
Receiver.stop(receiver)
Kernel.send(callback_handler, {:mammoth, :disconnected, true})
{:noreply, state}
end
def handle_cast(
:disconnected,
state = %{
subscriber: subscriber,
receiver: receiver,
callback_handler: callback_handler
}
) do
Subscriber.stop(subscriber)
Receiver.stop(receiver)
Kernel.send(callback_handler, {:mammoth, :disconnected, false})
{:noreply, 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