# SPDX-License-Identifier: MIT # Copyright (c) 2025 Scott Thompson defmodule ArduinoRouter.Socket do @moduledoc """ This is a transport module for the Arduino Router Bridge that uses the router's unix domain socket (a standard part of the Uno Q system) for communication. This method is not typically used directly. It's the defaut transport for the `ArduinoRouter.Bridge` and is typically created by that module. """ require Logger require :telemetry alias ArduinoRouter.MessagePackRPC # The default path to the Arduino router UNIX domain socket @default_socket_path ~c"/var/run/arduino-router.sock" @behaviour ArduinoRouter.Transport @typedoc """ Socket state structure. - `unix_socket`: The underlying Erlang socket handle - `buffer`: Accumulated bytes not yet parsed into complete messages """ @type t :: %__MODULE__{ unix_socket: :socket.socket() | nil, buffer: binary() } @typedoc "Result of connecting the socket to the arduino-router daemon" @type connect_result :: {:ok, t()} | {:error, term()} defstruct unix_socket: nil, buffer: <<>> @impl ArduinoRouter.Transport def start(receiver, opts \\ []) do socket_path = Keyword.get(opts, :socket_path, @default_socket_path) sockaddr = %{family: :local, path: socket_path} with {:ok, unix_socket} <- :socket.open(:local, :stream, :default), :ok <- :socket.connect(unix_socket, sockaddr) do :telemetry.execute([:arduino_router, :socket, :connected], %{}, %{ socket_path: socket_path }) socket = %__MODULE__{unix_socket: unix_socket, buffer: <<>>} run_receive_loop(socket, receiver) {:ok, socket} else {:error, reason} = error_result -> Logger.error("Failed to connect to Arduino router socket: #{inspect(reason)}") :telemetry.execute([:arduino_router, :socket, :connection_failed], %{}, %{ reason: reason }) error_result end end @impl ArduinoRouter.Transport def send_message(socket, message) when not is_nil(socket.unix_socket) do :telemetry.execute([:arduino_router, :socket, :send_message], %{}, %{message: message}) unpacked_message = MessagePackRPC.unpacked_message(message) with {:ok, packed_message} <- Msgpax.pack(unpacked_message), :ok <- :socket.send(socket.unix_socket, packed_message) do {:ok, socket} else error_result -> Logger.error( "Failed to pack and send message #{inspect(message)} to Arduino router: #{inspect(error_result)}" ) error_result end end def send_message(socket, _message) when is_nil(socket.unix_socket) do Logger.error("Cannot send message, socket not connected") {:error, :socket_not_connected} end @impl ArduinoRouter.Transport def stop(socket) do if socket.unix_socket do :telemetry.execute([:arduino_router, :socket, :closed], %{}, %{}) :socket.close(socket.unix_socket) end :ok end @spec run_receive_loop(t(), pid()) :: pid() defp run_receive_loop(socket, pid) do spawn_link(fn -> receive_loop(socket.unix_socket, pid, <<>>) end) end @spec receive_loop(:socket.socket(), pid(), binary()) :: no_return() defp receive_loop(socket, pid, buffer) when is_binary(buffer) do case send_messages(buffer, pid) do {:ok, rest} -> case :socket.recv(socket, 0) do {:ok, data} -> :telemetry.execute([:arduino_router, :socket, :received_data], %{}, %{data: data}) receive_loop(socket, pid, rest <> data) {:error, reason} -> Logger.warning("Socket connection closed: #{inspect(reason)}") exit({:socket_closed, reason}) end {:error, reason} -> Logger.error("Failed to process messages: #{inspect(reason)}") exit({:message_error, reason}) end end @spec send_messages(binary(), pid()) :: {:ok, binary()} | {:error, term()} defp send_messages(buffer, pid) do case Msgpax.unpack_slice(buffer) do {:ok, unpacked_message, rest} -> # Convert from an unpacked arry to an rpc message rpc_message = MessagePackRPC.to_rpc_message(unpacked_message) :telemetry.execute([:arduino_router, :socket, :incoming_rpc], %{}, %{message: rpc_message}) send(pid, {:incoming_rpc, rpc_message}) send_messages(rest, pid) {:error, %Msgpax.UnpackError{reason: :incomplete}} -> {:ok, buffer} {:error, error} = error_result -> Logger.error("Failed to unpack MessagePack-RPC message: #{inspect(error)}") error_result end end end