Packages

Drive R from Elixir through a persistent Rscript backend, with an experimental embedded native backend.

Current section

Files

Jump to
rx lib rx runtime.ex
Raw

lib/rx/runtime.ex

defmodule Rx.Runtime do
@moduledoc false
use GenServer
def start_link(opts \\ []) do
GenServer.start_link(__MODULE__, opts, name: __MODULE__)
end
def system_init(config) do
GenServer.call(__MODULE__, {:system_init, config}, :infinity)
|> handle_init_reply()
end
def use_backend(config) do
GenServer.call(__MODULE__, {:use_backend, config}, :infinity)
|> handle_init_reply()
end
def ensure_system_init(config) do
GenServer.call(__MODULE__, {:ensure_system_init, config}, :infinity)
|> handle_init_reply()
end
defp handle_init_reply(reply) do
case reply do
:ok ->
:ok
{:error, %_{} = exception} ->
raise exception
{:error, {:backend_not_switchable, message}} when is_binary(message) ->
raise RuntimeError, message
{:error, reason} ->
raise RuntimeError, "R backend initialization failed: #{inspect(reason)}"
end
end
def eval(source, globals, opts) do
case GenServer.call(__MODULE__, {:eval, source, globals, opts}, :infinity) do
{:ok, eval_payload} -> wrap_eval_result(eval_payload, opts)
{:error, :backend_crashed} -> raise RuntimeError, "R backend crashed"
{:error, :backend_not_running} -> raise RuntimeError, "R backend not running"
{:error, {:backend_exit, _}} -> raise RuntimeError, "R backend crashed"
{:error, %ArgumentError{} = error} -> raise error
{:error, %Rx.Error{} = error} -> raise error
{:error, error_map} when is_map(error_map) -> raise_r_error(error_map)
{:error, reason} -> raise RuntimeError, "R eval failed: #{inspect(reason)}"
end
end
def plot(source, globals, opts) do
case GenServer.call(__MODULE__, {:plot, source, globals, opts}, :infinity) do
{:ok, plot_payload} -> wrap_plot_result(plot_payload, opts)
{:error, :backend_crashed} -> raise RuntimeError, "R backend crashed"
{:error, :backend_not_running} -> raise RuntimeError, "R backend not running"
{:error, {:backend_exit, _}} -> raise RuntimeError, "R backend crashed"
{:error, %ArgumentError{} = error} -> raise error
{:error, %Rx.Error{} = error} -> raise error
{:error, error_map} when is_map(error_map) -> raise_r_error(error_map)
{:error, reason} -> raise RuntimeError, "R plot failed: #{inspect(reason)}"
end
end
def encode(type, value) do
case GenServer.call(__MODULE__, {:encode, type, value}, :infinity) do
{:ok, %Rx.Object{} = object} -> object
%Rx.Object{} = object -> object
{:error, reason} -> raise ArgumentError, "failed to encode R value: #{inspect(reason)}"
end
end
def print(%Rx.Object{} = object, opts) do
case GenServer.call(__MODULE__, {:print, object, opts}, :infinity) do
{:ok, output} -> wrap_print_result(output, opts)
{:error, :backend_crashed} -> raise RuntimeError, "R backend crashed"
{:error, :backend_not_running} -> raise RuntimeError, "R backend not running"
{:error, {:backend_exit, _}} -> raise RuntimeError, "R backend crashed"
{:error, %ArgumentError{} = error} -> raise error
{:error, %Rx.Error{} = error} -> raise error
{:error, error_map} when is_map(error_map) -> raise_r_error(error_map)
{:error, reason} -> raise RuntimeError, "R print failed: #{inspect(reason)}"
end
end
def decode(%Rx.Object{} = object) do
if inline_object?(object) do
decode_inline_object(object)
else
decode_backend_object(object)
end
end
defp decode_backend_object(%Rx.Object{} = object) do
case GenServer.call(__MODULE__, {:decode, object}, :infinity) do
{:ok, value} -> normalize_decoded(value)
{:opaque, object} -> object
{:error, reason} -> raise ArgumentError, "failed to decode R value: #{inspect(reason)}"
end
end
defp decode_inline_object(%Rx.Object{} = object) do
case Rx.Backends.PortArrow.decode(object) do
{:ok, value} ->
value
|> normalize_decoded()
|> restore_inline_source_values(object)
{:opaque, object} ->
object
{:error, reason} ->
raise ArgumentError, "failed to decode R value: #{inspect(reason)}"
end
end
defp restore_inline_source_values(decoded, %Rx.Object{
remote_info: {:inline_named_list, source_map}
})
when is_map(source_map) and not is_struct(source_map) do
restore_inline_source_value(decoded, source_map)
end
defp restore_inline_source_values(decoded, %Rx.Object{}), do: decoded
defp restore_inline_source_value(decoded, source_map)
when is_map(source_map) and not is_struct(source_map) do
Map.new(decoded, fn {key, value} ->
{key, restore_inline_source_value(value, Map.get(source_map, key))}
end)
end
defp restore_inline_source_value(
_decoded,
%Rx.Object{id: <<"__inline__:", _json::binary>>} = object
),
do: decode_inline_object(object)
defp restore_inline_source_value(_decoded, %Rx.Object{} = object), do: object
defp restore_inline_source_value(decoded, _source), do: decoded
defp inline_object?(%Rx.Object{id: "__null__"}), do: true
defp inline_object?(%Rx.Object{id: <<"__inline__:", _json::binary>>}), do: true
defp inline_object?(%Rx.Object{}), do: false
def decode_arrow(%Rx.Object{} = object) do
case GenServer.call(__MODULE__, {:decode_arrow, object}, :infinity) do
{:ok, binary} -> {:ok, binary}
{:error, reason} -> {:error, reason}
end
end
def decode_data_frame(%Rx.Object{} = object, opts) when is_list(opts) do
case GenServer.call(__MODULE__, {:decode_data_frame, object, opts}, :infinity) do
{:ok, wire} -> {:ok, wire}
{:error, reason} -> {:error, reason}
end
end
def encode_dataframe(ipc_bytes) when is_binary(ipc_bytes) do
case GenServer.call(__MODULE__, {:encode_dataframe, ipc_bytes}, :infinity) do
{:ok, %Rx.Object{} = obj} -> {:ok, obj}
{:error, reason} -> {:error, reason}
end
end
def encode_data_frame(wire, opts) when is_map(wire) and is_list(opts) do
case GenServer.call(__MODULE__, {:encode_data_frame, wire, opts}, :infinity) do
{:ok, %Rx.Object{} = obj} -> {:ok, obj}
{:error, reason} -> {:error, reason}
end
end
def backend do
GenServer.call(__MODULE__, :backend)
end
@impl true
def init(_opts) do
{:ok, %{init: nil, backend: nil}}
end
@impl true
def handle_call({:system_init, config}, _from, %{init: nil} = state) do
{reply, state} = init_backend(config, state)
{:reply, reply, state}
end
def handle_call({:use_backend, config}, _from, %{init: nil} = state) do
{reply, state} = init_backend(config, state)
{:reply, reply, state}
end
def handle_call({:use_backend, config}, _from, %{init: _existing} = state) do
if state.init == init_identity(config) do
{reply, state} = init_backend(config, state)
{:reply, reply, state}
else
case shutdown_backend(state.backend) do
:ok ->
{reply, state} = init_backend(config, %{state | init: nil, backend: nil})
{:reply, reply, state}
{:error, reason} ->
{:reply, {:error, reason}, state}
end
end
end
def handle_call({:ensure_system_init, config}, _from, %{init: nil} = state) do
configs = List.wrap(config)
{reply, state} = ensure_first_available(configs, state)
{:reply, reply, state}
end
def handle_call({:ensure_system_init, _config}, _from, %{init: existing} = state) do
# Re-run system_init if the backend process died since last init.
case safe_backend_call(fn -> state.backend.system_init(existing) end) do
:ok -> {:reply, :ok, state}
{:error, reason} -> {:reply, {:error, reason}, state}
end
end
def handle_call({:system_init, config}, _from, %{init: existing} = state) do
if existing == init_identity(config) do
# Re-run system_init if the backend process died since last init.
{reply, state} = init_backend(config, state)
{:reply, reply, state}
else
{:reply,
{:error, RuntimeError.exception("Rx is already initialized with different options")},
state}
end
end
def handle_call(:backend, _from, state) do
{:reply, backend_name(state.backend), state}
end
def handle_call({:eval, _source, _globals, _opts}, _from, %{backend: nil} = state) do
{:reply, {:error, :backend_not_running}, state}
end
def handle_call({:eval, source, globals, opts}, _from, state) do
case safe_backend_call(fn -> state.backend.eval(source, globals, opts) end) do
{:ok, payload} -> {:reply, {:ok, payload}, state}
{:error, :backend_crashed} -> {:reply, {:error, :backend_crashed}, state}
{:error, :backend_not_running} -> {:reply, {:error, :backend_not_running}, state}
{:error, {:backend_exit, _} = reason} -> {:reply, {:error, reason}, state}
{:error, reason} -> {:reply, {:error, reason}, state}
end
end
def handle_call({:plot, _source, _globals, _opts}, _from, %{backend: nil} = state) do
{:reply, {:error, :backend_not_running}, state}
end
def handle_call({:plot, source, globals, opts}, _from, state) do
case safe_backend_call(fn -> state.backend.plot(source, globals, opts) end) do
{:ok, payload} -> {:reply, {:ok, payload}, state}
{:error, :backend_crashed} -> {:reply, {:error, :backend_crashed}, state}
{:error, :backend_not_running} -> {:reply, {:error, :backend_not_running}, state}
{:error, {:backend_exit, _} = reason} -> {:reply, {:error, reason}, state}
{:error, reason} -> {:reply, {:error, reason}, state}
end
end
def handle_call({:encode, _type, _value}, _from, %{backend: nil} = state) do
{:reply, {:error, :backend_not_running}, state}
end
def handle_call({:encode, type, value}, _from, state) do
{:reply, safe_backend_call(fn -> state.backend.encode(type, value) end), state}
end
def handle_call({:print, %Rx.Object{}, _opts}, _from, %{backend: nil} = state) do
{:reply, {:error, :backend_not_running}, state}
end
def handle_call({:print, %Rx.Object{} = object, opts}, _from, state) do
{:reply, safe_backend_call(fn -> state.backend.print(object, opts) end), state}
end
def handle_call({:decode, %Rx.Object{}}, _from, %{backend: nil} = state) do
{:reply, {:error, :backend_not_running}, state}
end
def handle_call({:decode, %Rx.Object{} = object}, _from, state) do
case safe_backend_call(fn -> state.backend.decode(object) end) do
{:ok, value} -> {:reply, {:ok, value}, state}
{:opaque, obj} -> {:reply, {:opaque, obj}, state}
{:error, reason} -> {:reply, {:error, reason}, state}
end
end
def handle_call({:decode_arrow, %Rx.Object{}}, _from, %{backend: nil} = state) do
{:reply, {:error, :backend_not_running}, state}
end
def handle_call({:decode_arrow, %Rx.Object{} = object}, _from, state) do
result = route_decode_arrow(object, state)
{:reply, result, state}
end
def handle_call({:decode_data_frame, %Rx.Object{}, _opts}, _from, %{backend: nil} = state) do
{:reply, {:error, :backend_not_running}, state}
end
def handle_call({:decode_data_frame, %Rx.Object{} = object, opts}, _from, state) do
result = route_decode_data_frame(object, opts, state)
{:reply, result, state}
end
def handle_call({:encode_dataframe, _ipc_bytes}, _from, %{backend: nil} = state) do
{:reply, {:error, :backend_not_running}, state}
end
def handle_call({:encode_dataframe, ipc_bytes}, _from, state) when is_binary(ipc_bytes) do
result = route_encode_dataframe(ipc_bytes, state)
{:reply, result, state}
end
def handle_call({:encode_data_frame, _wire, _opts}, _from, %{backend: nil} = state) do
{:reply, {:error, :backend_not_running}, state}
end
def handle_call({:encode_data_frame, wire, opts}, _from, state)
when is_map(wire) and is_list(opts) do
result = route_encode_data_frame(wire, opts, state)
{:reply, result, state}
end
defp route_decode_arrow(%Rx.Object{} = object, state) do
if native_backend?(state.backend) do
route_native_decode_arrow(object, state)
else
safe_backend_call(fn -> state.backend.decode_arrow(object) end)
end
end
defp route_native_decode_arrow(%Rx.Object{} = object, state) do
case object.backend do
:port_arrow ->
with :ok <- ensure_dataframe_process_backend(state) do
safe_backend_call(fn -> dataframe_process_backend().decode_arrow(object) end)
end
_other ->
safe_backend_call(fn -> state.backend.decode_arrow(object) end)
end
end
defp route_decode_data_frame(%Rx.Object{} = object, opts, state) do
if native_backend?(state.backend) do
route_native_decode_data_frame(object, opts, state)
else
safe_backend_call(fn -> state.backend.decode_data_frame(object, opts) end)
end
end
defp route_native_decode_data_frame(%Rx.Object{backend: :port_arrow} = object, opts, state) do
with :ok <- ensure_dataframe_process_backend(state) do
safe_backend_call(fn -> dataframe_process_backend().decode_data_frame(object, opts) end)
end
end
defp route_native_decode_data_frame(%Rx.Object{} = object, opts, state) do
safe_backend_call(fn -> state.backend.decode_data_frame(object, opts) end)
end
defp route_encode_dataframe(ipc_bytes, state) do
if native_backend?(state.backend) do
route_native_encode_dataframe(ipc_bytes, state)
else
safe_backend_call(fn -> state.backend.encode_dataframe(ipc_bytes) end)
end
end
defp route_native_encode_dataframe(ipc_bytes, state) do
safe_backend_call(fn -> state.backend.encode_dataframe(ipc_bytes) end)
end
defp route_encode_data_frame(wire, opts, state) do
safe_backend_call(fn -> state.backend.encode_data_frame(wire, opts) end)
end
defp ensure_dataframe_process_backend(state) do
config = dataframe_process_config(state)
safe_backend_call(fn -> dataframe_process_backend().system_init(config) end)
end
defp dataframe_process_config(%{init: init}) when is_list(init) do
config = [backend: dataframe_process_backend()]
case Keyword.fetch(init, :lib_paths) do
{:ok, lib_paths} -> Keyword.put(config, :lib_paths, lib_paths)
:error -> config
end
end
defp dataframe_process_config(_state), do: [backend: dataframe_process_backend()]
defp dataframe_process_backend do
Application.get_env(:rx, :dataframe_process_backend, Rx.Backends.PortArrow)
end
defp native_backend?(backend) do
backend == Application.get_env(:rx, :runtime_native_backend, Rx.Backends.Native)
end
defp ensure_first_available([], state), do: {{:error, :backend_not_running}, state}
defp ensure_first_available([config | rest], state) do
case init_backend(config, state) do
{:ok, state} ->
{:ok, state}
{{:error, reason}, _state} when rest != [] ->
if continue_init_fallback?(config, reason) do
ensure_first_available(rest, state)
else
{{:error, reason}, state}
end
{{:error, reason}, _state} ->
{{:error, reason}, state}
end
end
defp continue_init_fallback?(config, reason) do
backend = Keyword.fetch!(config, :backend)
if native_backend?(backend) do
safe_native_init_fallback_reason?(reason)
else
true
end
end
defp safe_native_init_fallback_reason?({:embedded_nif_unavailable, _message}), do: true
defp safe_native_init_fallback_reason?({:not_loaded, _message}), do: true
defp safe_native_init_fallback_reason?(:missing_r_home), do: true
defp safe_native_init_fallback_reason?(:missing_lib_r), do: true
defp safe_native_init_fallback_reason?({:native_init_failed, %{retryable: true}}), do: true
defp safe_native_init_fallback_reason?({:native_init_failed, diagnostics})
when is_map(diagnostics),
do: false
defp safe_native_init_fallback_reason?({:native_init_mismatch, _diagnostics}), do: false
defp safe_native_init_fallback_reason?(_reason), do: false
defp init_backend(config, state) do
backend = Keyword.fetch!(config, :backend)
case safe_backend_call(fn -> backend.system_init(config) end) do
:ok ->
{:ok, %{state | init: init_identity(config), backend: backend}}
{:error, reason} ->
{{:error, reason}, state}
end
end
defp shutdown_backend(nil), do: :ok
defp shutdown_backend(backend) do
if function_exported?(backend, :shutdown, 0) do
safe_backend_call(fn -> backend.shutdown() end)
else
{:error,
{:backend_not_switchable,
"#{inspect(backend)} cannot be switched because it does not expose shutdown/0"}}
end
end
defp safe_backend_call(fun) do
fun.()
rescue
error in [ArgumentError, RuntimeError, Rx.Error] -> {:error, error}
catch
:exit, {:noproc, _} -> {:error, :backend_not_running}
:exit, {:normal, _} -> {:error, :backend_not_running}
:exit, reason -> {:error, {:backend_exit, reason}}
end
defp backend_name(nil), do: nil
defp backend_name(Rx.Backends.PortArrow), do: :port_arrow
defp backend_name(Rx.Backends.Native), do: :native
defp backend_name(backend), do: backend
defp init_identity(config) when is_list(config) do
Keyword.update(config, :renv, nil, &renv_identity/1)
end
defp init_identity(config), do: config
defp renv_identity(nil), do: nil
defp renv_identity(%{identity_env: identity_env} = renv) do
%{renv | env: identity_env}
end
defp renv_identity(%{env: env} = renv) when is_list(env) do
%{renv | env: Enum.sort_by(env, fn {name, value} -> {name, value} end)}
end
defp renv_identity(renv), do: renv
defp wrap_eval_result({result, globals, output}, opts) do
if opts[:capture] do
%Rx.EvalResult{
result: result,
globals: globals,
stdout: output.stdout,
messages: output.messages,
warnings: output.warnings
}
else
write_output!(output, opts)
{result, globals}
end
end
defp wrap_print_result(output, opts) do
if opts[:capture] do
%Rx.PrintResult{
stdout: output.stdout,
messages: output.messages,
warnings: output.warnings
}
else
output.stdout
end
end
defp wrap_plot_result(%{plots: plots, output: output}, opts) do
if opts[:capture] do
%Rx.PlotResult{
plots: plots,
stdout: output.stdout,
messages: output.messages,
warnings: output.warnings
}
else
write_output!(output, opts)
plots
end
end
defp write_output!(output, opts) do
unless output.stdout == "", do: IO.write(opts[:stdout_device], output.stdout)
stderr = opts[:stderr_device]
unless output.messages == "", do: IO.write(stderr, output.messages)
unless output.warnings == "", do: IO.write(stderr, output.warnings)
end
defp raise_r_error(%{message: message} = error_map) do
raise Rx.Error,
message: message,
r_class: normalize_r_character_vector(Map.get(error_map, :r_class)),
call: normalize_r_character_scalar(Map.get(error_map, :call)),
traceback: normalize_r_character_vector(Map.get(error_map, :traceback)),
output: Map.get(error_map, :output)
end
defp normalize_r_character_scalar(nil), do: nil
defp normalize_r_character_scalar([value | _rest]) when is_binary(value), do: value
defp normalize_r_character_scalar(value) when is_binary(value), do: value
defp normalize_r_character_scalar(value), do: inspect(value)
defp normalize_r_character_vector(nil), do: nil
defp normalize_r_character_vector(values) when is_list(values), do: values
defp normalize_r_character_vector(value) when is_binary(value), do: [value]
defp normalize_r_character_vector(value), do: [inspect(value)]
defp normalize_decoded({:na, type}), do: %Rx.NA{type: type}
defp normalize_decoded(values) when is_list(values), do: Enum.map(values, &normalize_decoded/1)
defp normalize_decoded(%Rx.RList{items: items} = rlist) do
%{rlist | items: Enum.map(items, fn {k, v} -> {k, normalize_decoded(v)} end)}
end
defp normalize_decoded(value) when is_map(value) and not is_struct(value) do
Map.new(value, fn {k, v} -> {k, normalize_decoded(v)} end)
end
defp normalize_decoded(value), do: value
end