Current section

Files

Jump to
urbit_ex lib api worker.ex
Raw

lib/api/worker.ex

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