defmodule UrbitEx.Server do alias UrbitEx.{API} use GenServer require Logger @moduledoc """ Documentation for `Urbitex`. """ def start_link(options) when is_list(options) do GenServer.start_link(__MODULE__, options, name: :urbit) end # callbacks @impl true def init(args) do [{:url, url}, {:code, code}] = args session = API.init(url, code) |> API.login() |> API.open_channel() |> API.start_sse() {:ok, session} end @impl true def handle_info(%HTTPoison.AsyncStatus{}, session) do {:noreply, session} end @impl true # ignore keep-alive messages def handle_info(%{chunk: "\n"}, session) do {:noreply, session} end @impl true def handle_info(%{chunk: data}, session) do {event_id, message} = data |> parse_stream new_session = %{session | last_sse: event_id, events: [message | session.events]} broadcast(session, message) {:noreply, new_session} end @impl true def handle_info(message, session) do IO.inspect(message, label: :handle_info_message) {:noreply, session} end defp parse_stream(event) do [_, id, _, data] = event |> String.split("\n") |> Enum.map(fn x -> String.split(x, ": ") end) |> List.flatten() |> Enum.filter(fn x -> String.length(x) > 0 end) {String.to_integer(id), Jason.decode!(data)} end def handle_call(:get, _from, session) do {:reply, session, session} end @impl true def handle_cast({:subscribe, subscription}, session) do session = API.subscribe(session, subscription) session = %{session | subscriptions: [subscription | session.subscriptions]} {:noreply, session} end @impl true def handle_cast({:consume, pid}, session) do {:noreply, %{session | consumers: [pid | session.consumers]}} end def broadcast(session, message) do session.consumers |> Enum.each(&send(&1, message)) end end