Packages
grizzly
0.4.3
9.1.4
9.1.2
9.1.1
9.1.0
9.0.0
8.15.3
8.15.2
8.15.1
8.15.0
8.14.0
8.13.0
8.12.0
8.11.3
8.11.2
8.11.1
8.11.0
8.10.0
8.9.0
8.8.1
8.8.0
8.7.1
8.7.0
8.6.12
8.6.11
8.6.10
8.6.9
8.6.8
8.6.7
retired
8.6.6
8.6.5
8.6.4
8.6.3
8.6.2
8.6.1
8.6.0
8.5.3
8.5.2
8.5.1
8.5.0
8.4.0
8.3.0
8.2.3
8.2.2
8.2.1
8.2.0
8.1.0
8.0.1
8.0.0
7.4.3
7.4.2
7.4.1
7.4.0
7.3.0
7.2.0
7.1.4
7.1.3
7.1.2
7.1.1
7.1.0
7.0.4
7.0.3
7.0.2
7.0.1
7.0.0
6.8.8
6.8.7
6.8.6
6.8.5
6.8.4
6.8.3
6.8.2
6.8.1
6.8.0
6.7.1
6.7.0
6.6.1
6.6.0
6.5.1
6.5.0
6.4.0
6.3.0
6.2.0
6.1.1
6.1.0
6.0.1
6.0.0
5.4.1
5.4.0
5.3.0
5.2.8
5.2.7
5.2.6
5.2.5
5.2.4
5.2.3
5.2.2
5.2.1
5.2.0
5.1.2
5.1.1
5.1.0
5.0.2
5.0.1
5.0.0
4.0.1
4.0.0
3.0.0
2.1.0
2.0.0
1.0.1
1.0.0
0.22.7
0.22.6
0.22.5
0.22.4
0.22.3
0.22.2
0.22.1
0.22.0
0.21.1
0.21.0
0.20.2
0.20.1
0.20.0
0.19.1
0.19.0
0.18.3
0.18.2
0.18.1
0.18.0
0.17.7
0.17.6
0.17.5
0.17.4
0.17.3
0.17.2
0.17.1
0.17.0
0.16.2
0.16.1
0.16.0
0.15.11
0.15.10
0.15.9
0.15.8
0.15.7
0.15.6
0.15.5
0.15.4
0.15.3
0.15.2
0.15.1
0.15.0
0.14.8
0.14.7
0.14.6
0.14.5
0.14.4
0.14.3
0.14.2
0.14.1
0.14.0
0.13.0
0.12.3
0.12.2
0.12.1
0.12.0
0.11.0
0.10.3
0.10.2
0.10.1
0.10.0
0.9.0
0.9.0-rc.4
0.9.0-rc.3
0.9.0-rc.2
0.9.0-rc.1
0.9.0-rc.0
0.8.8
0.8.7
0.8.6
0.8.5
0.8.4
0.8.3
0.8.2
0.8.1
0.8.0
0.7.0
0.6.6
0.6.5
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
0.5.0
0.4.3
0.4.2
Elixir Z-Wave library
Current section
Files
Jump to
Current section
Files
lib/grizzly/conn/server.ex
defmodule Grizzly.Conn.Server do
@moduledoc false
use GenServer
require Logger
alias Grizzly.Packet
alias Grizzly.{Notifications, Command, Conn}
alias Grizzly.Conn.Config
@retry_connect_delay 1_000
@type t :: pid
defmodule State do
@moduledoc false
@type command :: %{
from: GenServer.from(),
# the owner of queued commands, to which responses are to be sent
owner: pid,
command: Command.t(),
mode: Config.mode(),
status: :active | :queued,
queued_ref: nil | reference()
}
@type t :: %__MODULE__{
connected?: boolean,
socket: :inet.socket(),
config: Config.t(),
heart_beat_interval: pos_integer,
commands: [command]
}
defstruct connected?: false,
socket: nil,
config: nil,
heart_beat_interval: nil,
commands: []
def build_command(command, from, mode, opts) do
%{
command: command,
from: from,
owner: opts[:owner],
mode: mode,
status: :active,
queued_ref: nil
}
end
end
def child_spec(args) do
%{id: __MODULE__, start: {__MODULE__, :start_link, args}}
end
@spec start_link(Config.t()) :: GenServer.on_start()
def start_link(config) do
GenServer.start_link(__MODULE__, config)
end
@doc """
Check to see if the connection has been establish
"""
@spec connected?(pid) :: boolean
def connected?(conn) do
GenServer.call(conn, :connected?)
end
@spec send_command(Conn.t(), Command.t(), Keyword.t()) :: :ok | {:ok, any} | {:error, any}
def send_command(%Conn{conn: conn_server, mode: mode}, command, opts) do
GenServer.call(conn_server, {:send_command, command, mode, opts}, 120_000)
end
@doc """
Close the connection
"""
def close(conn) do
GenServer.call(conn, :close)
end
@impl true
def init(config) do
Kernel.send(self(), :setup)
{:ok, %State{config: config}}
end
@impl true
def handle_call(:connected?, _from, %State{socket: nil} = state) do
{:reply, false, state}
end
def handle_call(:connected?, _from, %State{socket: _socket} = state) do
{:reply, true, state}
end
def handle_call(:close, _, %State{socket: socket, config: config} = state) do
:ok = apply(config.client, :close, [socket])
{:reply, :ok, state}
end
def handle_call(
{:send_command, command, :sync, opts},
from,
%State{commands: commands} = state
) do
:ok = do_send_command(command, state)
command = State.build_command(command, from, :sync, opts)
{:noreply, %{state | commands: commands ++ [command]}}
end
def handle_call(
{:send_command, command, :async, opts},
from,
%State{commands: commands} = state
) do
:ok = do_send_command(command, state)
command = State.build_command(command, from, :async, opts)
{:reply, :ok, %{state | commands: commands ++ [command]}}
end
@impl true
def handle_info(:setup, %State{config: config} = state) do
case maybe_autoconnect(config) do
{:ok, socket} ->
_ = Logger.info("connected to: #{inspect(config.ip, base: :hex)}")
heart_beat_timer = heart_beat(config)
{:noreply, %{state | socket: socket, heart_beat_interval: heart_beat_timer}}
:noop ->
{:noreply, state}
{:error, :timeout} ->
_ =
Logger.warn(
"[GRIZZLY] Setup autoconnect timed out. Retrying in #{@retry_connect_delay}"
)
Process.send_after(self(), :setup, @retry_connect_delay)
{:noreply, state}
end
end
def handle_info(:heart_beat, %State{config: config, socket: socket} = state) do
apply(config.client, :send_heart_beat, [socket, [port: config.port]])
heart_beat_timer = heart_beat(config)
{:noreply, %{state | heart_beat_interval: heart_beat_timer}}
end
def handle_info(
data,
%State{config: config, commands: commands, connected?: connected?} = state
) do
case apply(config.client, :parse_response, [data]) do
{:ok, :heart_beat} ->
if !connected? do
Notifications.broadcast(:connection_established, config)
{:noreply, %{state | connected?: true}}
else
{:noreply, state}
end
{:ok, packet} ->
if Packet.ack_request?(packet) do
do_send_raw(Packet.as_ack_response(packet.seq_number), state)
{:noreply, state}
else
commands = process_commands(commands, packet, state)
{:noreply, %{state | commands: commands}}
end
{:error, :socket_closed} ->
_ = Logger.info("[GATEWAY]: Socket closed reconnecting")
_ = Process.cancel_timer(state.heart_beat_timer)
Kernel.send(self(), :setup)
{:noreply, %{state | connected?: false, heart_beat_timer: nil}}
end
end
defp heart_beat(%Config{heart_beat_timer: timer}) do
Process.send_after(self(), :heart_beat, timer)
end
defp maybe_autoconnect(%Config{autoconnect: false}), do: :noop
defp maybe_autoconnect(%Config{client: client, ip: ip, port: port}) do
_ = Logger.info("Attempting to connect to: #{inspect(ip, base: :hex)}")
apply(client, :connect, [ip, port])
end
defp do_send_command(command, %State{config: config, socket: socket}) do
client_opts = [ip_address: config.ip, port: config.port]
apply(config.client, :send, [socket, Command.encode(command), client_opts])
end
defp do_send_raw(binary, %State{config: config, socket: socket}) do
client_opts = [ip_address: config.ip, port: config.port]
apply(config.client, :send, [socket, binary, client_opts])
end
defp process_commands(commands, packet, state) do
commands
|> Enum.reduce(
[],
fn %{from: {pid, _ref} = sender, command: command, mode: mode} = cmd, acc ->
# Process managing the command may be gone (e.g. after a timeout)
if Process.alive?(command) do
case Command.handle_response(command, packet) do
{:finished, response} ->
send_response(cmd, response)
:ok = Command.complete(command)
acc
{:send_message, message} ->
case mode do
:sync ->
GenServer.reply(sender, message)
:async ->
send(pid, {:async_command, message})
end
acc ++ [cmd]
:retry ->
do_send_command(command, state)
acc ++ [cmd]
:continue ->
acc ++ [cmd]
:queued ->
new_cmd = handle_queued(cmd)
acc ++ [new_cmd]
end
else
acc
end
end
)
end
defp handle_queued(%{status: :queued} = cmd), do: cmd
defp handle_queued(%{status: :active, mode: :sync, from: from} = cmd) do
ref = make_ref()
:ok = GenServer.reply(from, {:ok, :queued, ref})
%{cmd | status: :queued, queued_ref: ref}
end
defp handle_queued(%{status: :active, mode: :aync, from: {pid, _}} = cmd) do
ref = make_ref()
:ok = send(pid, {:ok, :queued, ref})
%{cmd | status: :queued, queued_ref: ref}
end
defp send_response(
%{status: :queued, queued_ref: ref, from: {pid, _}, owner: owner},
response
) do
message = {Grizzly, :queued_response, ref, response}
receiver = owner || pid
_ =
Logger.info("[GRIZZLY] Sending queued response #{inspect(message)} to #{inspect(receiver)}")
send(receiver, {Grizzly, :queued_response, ref, response})
end
defp send_response(%{status: :active, from: {pid, _}, mode: :async}, response) do
send(pid, {:async_command, response})
end
defp send_response(%{status: :active, from: from, mode: :sync}, response) do
GenServer.reply(from, response)
end
end