Current section

Files

Jump to
grizzly lib grizzly connections async_connection.ex
Raw

lib/grizzly/connections/async_connection.ex

defmodule Grizzly.Connections.AsyncConnection do
@moduledoc false
# A connection type that is useful for doing long running operations that are
# allowed to be canceled, or if the operation may need to request more
# information from a user.
# don't use this connection type unless it is for a special reason. Normally,
# you will want to wrap this connection in a GenServer as normally there is
# some long running state tied to needing one of these connection types.
use GenServer
alias Grizzly.{Connections, Connection, Options, Report, Transport, ZIPGateway, ZWave}
alias Grizzly.Commands.CommandRunner
alias Grizzly.Connections.{KeepAlive, CommandList}
alias Grizzly.ZWave.Command
alias Grizzly.ZWave.Commands.ZIPPacket
require Logger
defmodule State do
@moduledoc false
defstruct transport: nil,
socket: nil,
commands: CommandList.empty(),
owner: nil,
keep_alive: nil,
node_id: nil
end
def child_spec(node_id, opts \\ []) do
%{id: __MODULE__, start: {__MODULE__, :start_link, [node_id, opts]}, restart: :transient}
end
@spec start_link(Options.t(), ZWave.node_id(), [Connection.opt()]) :: GenServer.on_start()
def start_link(grizzly_options, node_id, opts \\ []) do
name = Connections.make_name({:async, node_id})
opts = Keyword.put_new(opts, :owner, self())
GenServer.start_link(__MODULE__, [grizzly_options, node_id, opts], name: name)
end
@spec send_command(Grizzly.node_id(), Command.t(), keyword()) :: {:ok, reference()}
def send_command(node_id, command, opts \\ []) do
name = Connections.make_name({:async, node_id})
GenServer.call(name, {:send_command, command, opts}, 140_000)
end
@spec stop_command(Grizzly.node_id(), reference()) :: :ok
def stop_command(node_id, command_ref) do
name = Connections.make_name({:async, node_id})
GenServer.call(name, {:stop_command, command_ref})
end
@spec command_alive?(Grizzly.node_id(), reference()) :: boolean()
def command_alive?(node_id, command_ref) do
name = Connections.make_name({:async, node_id})
GenServer.call(name, {:command_alive?, command_ref})
end
def stop(node_id) do
# TODO close socket
name = Connections.make_name({:async, node_id})
GenServer.stop(name, :normal)
end
@impl GenServer
def init([grizzly_options, node_id, opts]) do
host = ZIPGateway.host_for_node(node_id, grizzly_options)
transport_impl = grizzly_options.transport
transport_opts = [
ip_address: host,
port: grizzly_options.zipgateway_port
]
case Transport.open(transport_impl, transport_opts) do
{:ok, transport} ->
{:ok,
%State{
transport: transport,
keep_alive: KeepAlive.init(node_id, 25_000),
owner: Keyword.fetch!(opts, :owner),
node_id: node_id
}}
{:error, reason} ->
{:stop, reason}
end
end
@impl GenServer
def handle_call({:send_command, command, send_opts}, {waiter, _ref}, state) do
{:ok, command_runner, command_ref, new_command_list} =
CommandList.create(state.commands, command, state.node_id, waiter, send_opts)
case do_send_command(command_runner, state) do
:ok ->
{:reply, {:ok, command_ref},
%State{
state
| commands: new_command_list,
keep_alive: KeepAlive.timer_restart(state.keep_alive)
}}
end
end
def handle_call({:stop_command, command_ref}, _from, state) do
{:ok, new_commands} = CommandList.stop_command_by_ref(state.commands, command_ref)
{:reply, :ok, %State{state | commands: new_commands}}
end
def handle_call({:command_alive?, command_ref}, _from, state) do
{:reply, CommandList.has_command_ref?(state.commands, command_ref), state}
end
@impl GenServer
def handle_info(:keep_alive_tick, state) do
%State{keep_alive: keep_alive} = state
new_keep_alive =
keep_alive
|> KeepAlive.make_command()
|> KeepAlive.run(&do_send_command(&1, state))
{:noreply, %State{state | keep_alive: new_keep_alive}}
end
# handle when there is a timeout and command runner stops
def handle_info(
{:grizzly, :command_timeout, command_runner_pid, grizzly_command},
state
) do
if grizzly_command.source.name == :keep_alive do
{:noreply, state}
else
waiter = CommandList.get_waiter_for_runner(state.commands, command_runner_pid)
do_timeout_reply(waiter, grizzly_command)
{:noreply,
%State{
state
| commands: CommandList.drop_command_runner(state.commands, command_runner_pid)
}}
end
end
def handle_info(data, state) do
%State{transport: transport, node_id: node_id} = state
case Transport.parse_response(transport, data) do
{:ok, :connection_closed} ->
Logger.debug("[Grizzly] connection to node #{inspect(node_id)} closed")
{:stop, :normal, state}
{:ok, transport_response} ->
updated_state = handle_commands(transport_response.command, state)
{:noreply, updated_state}
{:error, error} ->
error_message = Exception.message(error)
Logger.warn("[Grizzly] #{inspect(error_message)}")
{:noreply, state}
end
end
defp handle_commands(%Command{name: :keep_alive}, state) do
%State{state | keep_alive: KeepAlive.timer_restart(state.keep_alive)}
end
defp handle_commands(zip_packet, state) do
updated_state =
case CommandList.response_for_zip_packet(state.commands, zip_packet) do
{:retry, command_runner, new_command_list} ->
:ok = do_send_command(command_runner, state)
%State{state | commands: new_command_list}
{:continue, new_command_list} ->
if !ZIPPacket.ack_response?(zip_packet) do
# Since we are doing async communications we need to handle when the
# connection gets an unhandled command from the Z-Wave
send(
state.owner,
{:grizzly, :report, to_report(zip_packet, state.node_id)}
)
end
%State{state | commands: new_command_list}
{waiter, {:error, :nack_response, new_command_list}} ->
send(waiter, {:error, :nack_response})
%State{state | commands: new_command_list}
{waiter, {%Report{} = report, new_command_list}} ->
send(waiter, {:grizzly, :report, report})
%State{state | commands: new_command_list}
end
%State{updated_state | keep_alive: KeepAlive.timer_restart(state.keep_alive)}
end
defp to_report(zip_packet, node_id) do
Report.new(:complete, :command, node_id, command: Command.param!(zip_packet, :command))
end
defp do_send_command(command_runner, state) do
%State{transport: transport} = state
binary = CommandRunner.encode_command(command_runner)
Transport.send(transport, binary)
end
defp do_timeout_reply(waiter, grizzly_command) do
if grizzly_command.status == :queued do
{pid, _tag} = waiter
report =
Report.new(:complete, :timeout, grizzly_command.node_id,
command_ref: grizzly_command.ref,
queued: true
)
send(pid, {:grizzly, :report, report})
else
report =
Report.new(:complete, :timeout, grizzly_command.node_id, command_ref: grizzly_command.ref)
send(waiter, {:grizzly, :report, report})
end
end
end