Packages

A real-time time series database - command line client.

Retired package: I'll never finish it

Current section

Files

Jump to
phasedb_client lib phasedb client connection.ex
Raw

lib/phasedb/client/connection.ex

defmodule PhaseDB.Client.Connection do
alias PhaseDB.Client.ConnectionSupervisor
alias PhaseDB.Client.Request
alias PhaseDB.Result
require Logger
@behaviour :websocket_client
@moduledoc """
Implements the `websocket_client` behaviour and represents an active
connection to a phaseDB server.
"""
def create do
create PhaseDB.Client.default_uri
end
def create uri do
Supervisor.start_child ConnectionSupervisor, [uri]
end
def find do
find PhaseDB.Client.default_uri
end
def find uri do
:gproc.lookup_pids {:p, :l, {__MODULE__, uri}}
end
def find_or_create do
find_or_create PhaseDB.Client.default_uri
end
def find_or_create uri do
case find uri do
[] -> create uri
[pid | _] -> {:ok, pid}
end
end
def drop pid do
Supervisor.terminate_child ConnectionSupervisor, pid
end
def xmit pid, frame do
send pid, {:xmit, frame}
end
def start_link uri do
:websocket_client.start_link String.to_char_list(uri), __MODULE__, [uri]
end
def init [uri] do
:gproc.reg {:p, :l, {__MODULE__, uri}}
{:once, %{uri: uri, rx_frames: 0, tx_frames: 0, server_version: nil, engine_version: nil}}
end
def onconnect _req, state do
{:ok, state}
end
def ondisconnect {:remote, :closed}, state do
{:reconnect, %{state | rx_frames: 0, tx_frames: 0}}
end
def ondisconnect {:error, :econnrefused}, %{uri: uri} do
{:close, :error, "Connection Refused: #{uri}"}
end
def websocket_handle {:pong,_}, _req, state do
{:ok, state}
end
def websocket_handle {:text, capabilities}, _req, %{rx_frames: 0}=state do
case Result.from_json capabilities do
%Result{status: :ok}=response ->
response =
response
|> Enum.to_list
|> List.first
sv = Map.get response, :server_version
ev = Map.get response, :engine_version
Logger.info "Connected to PhaseDB Server #{sv}, Engine #{ev}"
{:ok, %{state | rx_frames: 1, server_version: sv, engine_version: ev}}
{:error, reason} ->
{:close, :error, reason}
end
end
def websocket_handle {:text, _}=frame, _req, %{rx_frames: rx_frames}=state do
:gproc.send {:p, :l, {Request, self}}, frame
{:ok, %{state | rx_frames: rx_frames+1}}
end
def websocket_info {:xmit, frame}, _req, %{tx_frames: tx_frames}=state do
{:reply, frame, %{state | tx_frames: tx_frames+1}}
end
def websocket_terminate _reason, _req, _state do
:ok
end
end