Packages
kujira
0.1.52
0.1.80
0.1.79
0.1.78
0.1.77
0.1.76
0.1.75
0.1.74
0.1.73
0.1.72
0.1.71
0.1.70
0.1.69
0.1.68
0.1.67
0.1.66
0.1.65
0.1.64
0.1.63
0.1.62
0.1.61
0.1.60
0.1.59
0.1.58
0.1.57
0.1.56
0.1.55
0.1.54
0.1.53
0.1.52
0.1.51
0.1.50
0.1.49
0.1.48
0.1.47
0.1.46
0.1.45
0.1.44
0.1.43
0.1.42
0.1.41
0.1.40
0.1.39
0.1.38
0.1.37
0.1.36
0.1.35
0.1.34
0.1.33
0.1.32
0.1.31
0.1.30
0.1.29
0.1.28
0.1.27
0.1.25
0.1.24
0.1.23
0.1.22
0.1.21
0.1.20
0.1.19
0.1.18
0.1.17
0.1.16
0.1.15
0.1.14
0.1.13
0.1.12
0.1.10
0.1.9
0.1.8
0.1.7
0.1.6
0.1.5
0.1.4
0.1.3
0.1.2
0.1.1
0.1.0
Elixir interfaces to Kujira dApps, for building indexers, APIs and bots
Current section
Files
Jump to
Current section
Files
lib/kujira/node/websocket.ex
defmodule Kujira.Node.Websocket do
defmacro __using__(opts) do
pubsub = Keyword.get(opts, :pubsub)
quote do
use WebSockex
require Logger
def start_link(config) do
endpoint = config[:websocket]
subscriptions = Keyword.get(config, :subscriptions, [])
Logger.info("#{__MODULE__} Starting node websocket: #{endpoint}")
{:ok, pid} = WebSockex.start_link("#{endpoint}/websocket", __MODULE__, %{})
for {s, idx} <- Enum.with_index(subscriptions), do: subscribe(pid, idx, s)
{:ok, pid}
end
def handle_connect(_conn, state) do
Logger.info("#{__MODULE__} Connected")
{:ok, state}
end
def handle_disconnect(status, _state) do
raise "#{__MODULE__} Disconnected"
end
def handle_frame({:text, msg}, state) do
case Jason.decode(msg, keys: :atoms) do
{:ok, %{id: id, result: %{data: %{type: t, value: v}}}} ->
Logger.info("#{__MODULE__} Subscription #{id} event #{t}")
Phoenix.PubSub.broadcast(unquote(pubsub), t, v)
{:ok, state}
{:ok, %{id: id, jsonrpc: "2.0", result: %{}}} ->
Logger.info("#{__MODULE__} Subscription #{id} successful")
{:ok, state}
_ ->
{:ok, state}
end
end
def handle_cast({:send, {_type, msg} = frame}, state) do
Logger.debug("#{__MODULE__} [send] #{msg}")
{:reply, frame, state}
end
defp subscribe(pid, id, query) do
message =
Jason.encode!(%{
jsonrpc: "2.0",
method: "subscribe",
id: id,
params: %{
query: query
}
})
WebSockex.send_frame(pid, {:text, message})
end
end
end
end