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