Current section
Files
Jump to
Current section
Files
lib/mavlink/local_connection.ex
defmodule XMAVLink.LocalConnection do
@moduledoc false
# XMAVLink.Router delegate for local connections, i.e
# Elixir processes using the Router API to subscribe to
# and send MAVLink messages.
require Logger
import Enum, only: [reduce: 3, filter: 2]
alias XMAVLink.Frame
alias XMAVLink.LocalConnection
defstruct system: nil,
component: nil,
subscriptions: [],
sequence_number: 0
@type t :: %LocalConnection{
system: 1..255,
component: 1..255,
subscriptions: [],
sequence_number: 0..255
}
# Handle message from Router.pack_and_send()
# We use handle_info instead of cast for symmetry
# with the other connection types
def handle_info(
{:local, frame},
receiving_connection = %LocalConnection{
system: system,
component: component,
sequence_number: sequence_number
},
_dialect
) do
# Fill in missing frame details source_system, source_component, sequence_number
{
:ok,
:local,
struct(receiving_connection, sequence_number: rem(sequence_number + 1, 255)),
struct(frame,
source_system: system,
source_component: component,
sequence_number: sequence_number
)
|> Frame.pack_frame()
}
end
def connect(:local, system, component) do
local_connection = struct(LocalConnection, system: system, component: component)
send(
# Local connection guaranteed, so this connect() called directly from Router process
self(),
{
:add_connection,
:local,
case Agent.start(fn -> [] end, name: XMAVLink.SubscriptionCache) do
{:ok, _} ->
:ok = Logger.debug("Started Subscription Cache")
# No subscriptions to restore
local_connection
{:error, {:already_started, _}} ->
:ok = Logger.debug("Restoring subscriptions from Subscription Cache")
reduce(
Agent.get(XMAVLink.SubscriptionCache, fn subs -> subs end),
local_connection,
fn {query, pid}, lc -> subscribe(query, pid, lc) end
)
end
}
)
end
def forward(to_connection, frame = %Frame{message: nil}) do
# If we couldn't unpack the message set the message_type to XMAVLink.UnknownMessage
forward(to_connection, struct(frame, message: %{__struct__: XMAVLink.UnknownMessage}))
end
def forward(
%LocalConnection{
subscriptions: subscriptions
},
frame = %Frame{
source_system: source_system,
source_component: source_component,
target_system: target_system,
target_component: target_component,
target: target,
message: message = %{__struct__: message_type}
}
) do
for {
%{
message: q_message_type,
source_system: q_source_system,
source_component: q_source_component,
target_system: q_target_system,
target_component: q_target_component,
as_frame: as_frame?
},
pid
} <- subscriptions do
if (q_message_type == nil or q_message_type == message_type) and
(q_source_system == 0 or q_source_system == source_system) and
(q_source_component == 0 or q_source_component == source_component) and
(q_target_system == 0 or
(target != :broadcast and target != :component and q_target_system == target_system)) and
(q_target_component == 0 or
(target != :broadcast and target != :system and
q_target_component == target_component)) do
send(pid, if(as_frame?, do: frame, else: message))
end
end
end
# Subscription request from subscriber
def subscribe(query, pid, local_connection) do
:ok = Logger.debug("Subscribe #{inspect(pid)} to query #{inspect(query)}")
# Monitor so that we can unsubscribe dead processes
Process.monitor(pid)
# Uniq prevents duplicate subscriptions
%LocalConnection{
local_connection
| subscriptions:
Enum.uniq([{query, pid} | local_connection.subscriptions])
|> update_subscription_cache
}
end
# Unsubscribe request from subscriber
def unsubscribe(pid, local_connection) do
:ok = Logger.debug("Unsubscribe #{inspect(pid)}")
%LocalConnection{
local_connection
| subscriptions:
filter(local_connection.subscriptions, &(not match?({_, ^pid}, &1)))
|> update_subscription_cache
}
end
# Automatically unsubscribe a dead subscriber process
def subscriber_down(pid, local_connection) do
:ok = Logger.debug("Subscriber #{inspect(pid)} exited")
%LocalConnection{
local_connection
| subscriptions:
filter(local_connection.subscriptions, &(not match?({_, ^pid}, &1)))
|> update_subscription_cache
}
end
defp update_subscription_cache(subscriptions) do
:ok = Logger.debug("Update subscription cache: #{inspect(subscriptions)}")
Agent.update(XMAVLink.SubscriptionCache, fn _ -> subscriptions end)
subscriptions
end
end