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