Page Menu
Home
Phorge
Search
Configure Global Search
Log In
Files
F85629156
No One
Temporary
Actions
View File
Edit File
Delete File
View Transforms
Subscribe
Award Token
Flag For Later
Size
28 KB
Referenced Files
None
Subscribers
None
View Options
diff --git a/c_src/exile.c b/c_src/exile.c
index adff507..47d5337 100644
--- a/c_src/exile.c
+++ b/c_src/exile.c
@@ -1,509 +1,509 @@
#include "erl_nif.h"
#include <errno.h>
#include <fcntl.h>
#include <signal.h>
#include <stdbool.h>
#include <stdio.h>
#include <string.h>
#include <sys/types.h>
#include <unistd.h>
//#define DEBUG
#ifdef DEBUG
#define debug(...) \
do { \
enif_fprintf(stderr, __VA_ARGS__); \
enif_fprintf(stderr, "\n"); \
} while (0)
#define start_timing() ErlNifTime __start = enif_monotonic_time(ERL_NIF_USEC)
#define elapsed_microseconds() (enif_monotonic_time(ERL_NIF_USEC) - __start)
#else
#define debug(...)
#define start_timing()
#define elapsed_microseconds() 0
#endif
#define error(...) \
do { \
enif_fprintf(stderr, __VA_ARGS__); \
enif_fprintf(stderr, "\n"); \
} while (0)
#define ERL_TRUE enif_make_atom(env, "true")
#define ERL_FALSE enif_make_atom(env, "false")
#define ERL_UNDEFINED enif_make_atom(env, "undefined")
#define MAKE_OK(term) enif_make_tuple2(env, ATOM_OK, term)
#define MAKE_ERROR(term) enif_make_tuple2(env, ATOM_ERROR, term)
#define GET_CTX(env, arg, ctx) \
do { \
ExilePriv *data = enif_priv_data(env); \
if (enif_get_resource(env, arg, data->rt, (void **)&ctx) == false) { \
return MAKE_ERROR(ATOM_INVALID_CTX); \
} \
} while (0);
static const int PIPE_READ = 0;
static const int PIPE_WRITE = 1;
static const int PIPE_CLOSED = -1;
static const int CMD_EXIT = -1;
static const int MAX_ARGUMENTS = 20;
static const int MAX_ARGUMENT_LEN = 1024;
static ERL_NIF_TERM ATOM_OK;
static ERL_NIF_TERM ATOM_ERROR;
static ERL_NIF_TERM ATOM_UNDEFINED;
static ERL_NIF_TERM ATOM_INVALID_CTX;
static ERL_NIF_TERM ATOM_PIPE_CLOSED;
enum exec_status {
SUCCESS,
PIPE_CREATE_ERROR,
PIPE_FLAG_ERROR,
FORK_ERROR,
PIPE_DUP_ERROR,
NULL_DEV_OPEN_ERROR,
};
typedef struct ExilePriv {
ErlNifResourceType *rt;
} ExilePriv;
typedef struct ExecContext {
int cmd_input_fd;
int cmd_output_fd;
int cmd_exit_status;
pid_t pid;
} ExecContext;
typedef struct ExecResult {
enum exec_status status;
int err;
ExecContext context;
} ExecResult;
static void rt_dtor(ErlNifEnv *env, void *obj) {
debug("Exile rt_dtor called\n");
}
static void rt_stop(ErlNifEnv *env, void *obj, int fd, int is_direct_call) {
debug("Exile rt_stop called\n");
}
static void rt_down(ErlNifEnv *env, void *obj, ErlNifPid *pid,
ErlNifMonitor *monitor) {
debug("Exile rt_down called\n");
}
static ErlNifResourceTypeInit rt_init = {rt_dtor, rt_stop, rt_down};
static int set_flag(int fd, int flags) {
return fcntl(fd, F_SETFL, fcntl(fd, F_GETFL) | flags);
}
static void close_all(int pipes[2][2]) {
for (int i = 0; i < 2; i++) {
if (pipes[i][PIPE_READ] > 0)
close(pipes[i][PIPE_READ]);
if (pipes[i][PIPE_WRITE] > 0)
close(pipes[i][PIPE_WRITE]);
}
}
#define RETURN_ERROR(__err) \
do { \
result.err = errno; \
result.status = __err; \
close_all(pipes); \
error("error in start_proccess(), %s:%d %s", __FILE__, __LINE__, \
strerror(errno)); \
return result; \
} while (0);
static ExecResult start_proccess(char *args[], bool stderr_to_console) {
ExecResult result;
pid_t pid;
int pipes[2][2] = {{0, 0}, {0, 0}};
if (pipe(pipes[STDIN_FILENO]) == -1 || pipe(pipes[STDOUT_FILENO]) == -1) {
debug("failed create pipes");
RETURN_ERROR(PIPE_CREATE_ERROR)
}
const int r_cmdin = pipes[STDIN_FILENO][PIPE_READ];
const int w_cmdin = pipes[STDIN_FILENO][PIPE_WRITE];
const int r_cmdout = pipes[STDOUT_FILENO][PIPE_READ];
const int w_cmdout = pipes[STDOUT_FILENO][PIPE_WRITE];
if (set_flag(r_cmdin, O_CLOEXEC) < 0 || set_flag(w_cmdout, O_CLOEXEC) < 0 ||
set_flag(w_cmdin, O_CLOEXEC | O_NONBLOCK) < 0 ||
set_flag(r_cmdout, O_CLOEXEC | O_NONBLOCK) < 0) {
debug("failed to set flag for pipes");
RETURN_ERROR(PIPE_FLAG_ERROR)
}
switch (pid = fork()) {
case -1:
debug("failed to fork");
RETURN_ERROR(FORK_ERROR)
case 0: // child
// close default stdio fd
close(STDIN_FILENO);
close(STDOUT_FILENO);
if (dup2(r_cmdin, STDIN_FILENO) < 0) {
debug("failed dup command input pipe to stdin");
RETURN_ERROR(PIPE_DUP_ERROR)
}
if (dup2(w_cmdout, STDOUT_FILENO) < 0) {
debug("failed dup command output pipe to stdout");
RETURN_ERROR(PIPE_DUP_ERROR)
}
if (stderr_to_console != true) {
close(STDERR_FILENO);
int dev_null = open("/dev/null", O_WRONLY);
if (dev_null == -1) {
RETURN_ERROR(NULL_DEV_OPEN_ERROR);
}
if (dup2(dev_null, STDERR_FILENO) < 0) {
debug("failed dup command error pipe to stderr");
close(dev_null);
RETURN_ERROR(PIPE_DUP_ERROR)
}
close(dev_null);
}
close_all(pipes);
execvp(args[0], args);
perror("execvp(): failed");
default: // parent
// close file descriptors used by child
close(r_cmdin);
close(w_cmdout);
result.context.pid = pid;
result.context.cmd_input_fd = w_cmdin;
result.context.cmd_output_fd = r_cmdout;
result.status = SUCCESS;
return result;
}
}
/* TODO: return appropriate error instead returning generic "badarg" error */
static ERL_NIF_TERM exec_proc(ErlNifEnv *env, int argc,
const ERL_NIF_TERM argv[]) {
char tmp[MAX_ARGUMENTS][MAX_ARGUMENT_LEN + 1];
char *exec_args[MAX_ARGUMENTS + 1];
unsigned int args_len;
if (enif_get_list_length(env, argv[0], &args_len) != true)
return enif_make_badarg(env);
if (args_len > MAX_ARGUMENTS)
return enif_make_badarg(env);
ERL_NIF_TERM head, tail, list = argv[0];
for (unsigned int i = 0; i < args_len; i++) {
if (enif_get_list_cell(env, list, &head, &tail) != true)
return enif_make_badarg(env);
if (enif_get_string(env, head, tmp[i], MAX_ARGUMENT_LEN, ERL_NIF_LATIN1) <
1)
return enif_make_badarg(env);
exec_args[i] = tmp[i];
list = tail;
}
exec_args[args_len] = NULL;
bool stderr_to_console = true;
int tmp_int;
if (enif_get_int(env, argv[1], &tmp_int) != true)
return enif_make_badarg(env);
stderr_to_console = tmp_int == 1 ? true : false;
struct ExilePriv *data = enif_priv_data(env);
ExecResult result = start_proccess(exec_args, stderr_to_console);
ExecContext *ctx = NULL;
ERL_NIF_TERM term;
switch (result.status) {
case SUCCESS:
ctx = enif_alloc_resource(data->rt, sizeof(ExecContext));
ctx->cmd_input_fd = result.context.cmd_input_fd;
ctx->cmd_output_fd = result.context.cmd_output_fd;
ctx->pid = result.context.pid;
debug("pid: %d cmd_in_fd: %d cmd_out_fd: %d", ctx->pid, ctx->cmd_input_fd,
ctx->cmd_output_fd);
term = enif_make_resource(env, ctx);
// resource should be collected beam GC when there are no more references
enif_release_resource(ctx);
return MAKE_OK(term);
default:
return MAKE_ERROR(enif_make_int(env, result.err));
}
}
static int select_write(ErlNifEnv *env, ExecContext *ctx) {
int retval = enif_select(env, ctx->cmd_input_fd, ERL_NIF_SELECT_WRITE, ctx,
NULL, ERL_UNDEFINED);
if (retval != 0)
perror("select_write()");
return retval;
}
static ERL_NIF_TERM write_proc(ErlNifEnv *env, int argc,
const ERL_NIF_TERM argv[]) {
if (argc != 2)
enif_make_badarg(env);
ExecContext *ctx = NULL;
GET_CTX(env, argv[0], ctx);
if (ctx->cmd_input_fd == PIPE_CLOSED)
return MAKE_ERROR(ATOM_PIPE_CLOSED);
ErlNifBinary bin;
if (enif_inspect_binary(env, argv[1], &bin) != true)
return enif_make_badarg(env);
unsigned int result = write(ctx->cmd_input_fd, bin.data, bin.size);
// TODO: cleanup
if (result >= bin.size) { // request completely satisfied
return MAKE_OK(enif_make_int(env, result));
} else if (result >= 0) { // request partially satisfied
int retval = select_write(env, ctx);
if (retval != 0)
return MAKE_ERROR(enif_make_int(env, retval));
return MAKE_OK(enif_make_int(env, result));
} else if (errno == EAGAIN) { // busy
int retval = select_write(env, ctx);
if (retval != 0)
return MAKE_ERROR(enif_make_int(env, retval));
return MAKE_ERROR(enif_make_int(env, EAGAIN));
} else { // Error
perror("write()");
return MAKE_ERROR(enif_make_int(env, errno));
}
}
static ERL_NIF_TERM close_pipe(ErlNifEnv *env, int argc,
const ERL_NIF_TERM argv[]) {
ExecContext *ctx = NULL;
GET_CTX(env, argv[0], ctx);
int kind;
enif_get_int(env, argv[1], &kind);
int result;
switch (kind) {
case 0:
if (ctx->cmd_input_fd == PIPE_CLOSED) {
return ATOM_OK;
} else {
result = close(ctx->cmd_input_fd);
if (result == 0) {
ctx->cmd_input_fd = PIPE_CLOSED;
return ATOM_OK;
} else {
perror("cmd_input_fd close()");
return MAKE_ERROR(enif_make_int(env, errno));
}
}
case 1:
if (ctx->cmd_output_fd == PIPE_CLOSED) {
return ATOM_OK;
} else {
result = close(ctx->cmd_output_fd);
if (result == 0) {
ctx->cmd_output_fd = PIPE_CLOSED;
return ATOM_OK;
} else {
perror("cmd_output_fd close()");
return MAKE_ERROR(enif_make_int(env, errno));
}
}
default:
debug("invalid file descriptor type");
return enif_make_badarg(env);
}
}
static int select_read(ErlNifEnv *env, ExecContext *ctx) {
int retval = enif_select(env, ctx->cmd_output_fd, ERL_NIF_SELECT_READ, ctx,
NULL, ERL_UNDEFINED);
if (retval != 0)
perror("select_read()");
return retval;
}
static ERL_NIF_TERM read_proc(ErlNifEnv *env, int argc,
const ERL_NIF_TERM argv[]) {
if (argc != 2)
enif_make_badarg(env);
ExecContext *ctx = NULL;
GET_CTX(env, argv[0], ctx);
if (ctx->cmd_output_fd == PIPE_CLOSED)
return MAKE_ERROR(ATOM_PIPE_CLOSED);
bool is_buffered = true;
int size;
enif_get_int(env, argv[1], &size);
if (size == -1) {
size = 65535;
is_buffered = false;
} else if (size > 65535 || size < 1) {
enif_make_badarg(env);
}
unsigned char buf[size];
int result = read(ctx->cmd_output_fd, buf, sizeof(buf));
ERL_NIF_TERM bin_term;
if (result >= 0) {
ErlNifBinary bin;
enif_alloc_binary(result, &bin);
// TODO: we should use binary when reading itself instead of allocating
// again
memcpy(bin.data, buf, result);
bin_term = enif_make_binary(env, &bin);
}
// TODO: cleanup
if (result >= size ||
(is_buffered == false && result >= 0)) { // request completely satisfied
return MAKE_OK(bin_term);
} else if (result > 0) { // request partially satisfied
int retval = select_read(env, ctx);
if (retval != 0)
return MAKE_ERROR(enif_make_int(env, retval));
return MAKE_OK(bin_term);
} else if (result == 0) { // EOF
return MAKE_OK(bin_term);
} else if (errno == EAGAIN) { // busy
int retval = select_read(env, ctx);
if (retval != 0)
return MAKE_ERROR(enif_make_int(env, retval));
return MAKE_ERROR(enif_make_int(env, EAGAIN));
} else { // Error
perror("read()");
return MAKE_ERROR(enif_make_int(env, errno));
}
}
static ERL_NIF_TERM is_alive(ErlNifEnv *env, int argc,
const ERL_NIF_TERM argv[]) {
ExecContext *ctx = NULL;
GET_CTX(env, argv[0], ctx);
if (ctx->pid == CMD_EXIT)
return ERL_FALSE;
int result = kill(ctx->pid, 0);
if (result == 0) {
return ERL_TRUE;
} else {
return ERL_FALSE;
}
}
static ERL_NIF_TERM terminate_proc(ErlNifEnv *env, int argc,
const ERL_NIF_TERM argv[]) {
ExecContext *ctx = NULL;
GET_CTX(env, argv[0], ctx);
if (ctx->pid == CMD_EXIT)
return MAKE_OK(enif_make_int(env, 0));
return enif_make_int(env, kill(ctx->pid, SIGTERM));
}
static ERL_NIF_TERM kill_proc(ErlNifEnv *env, int argc,
const ERL_NIF_TERM argv[]) {
ExecContext *ctx = NULL;
GET_CTX(env, argv[0], ctx);
if (ctx->pid == CMD_EXIT)
return MAKE_OK(enif_make_int(env, 0));
return enif_make_int(env, kill(ctx->pid, SIGKILL));
}
static ERL_NIF_TERM wait_proc(ErlNifEnv *env, int argc,
const ERL_NIF_TERM argv[]) {
ExecContext *ctx = NULL;
GET_CTX(env, argv[0], ctx);
if (ctx->pid == CMD_EXIT)
return MAKE_OK(enif_make_int(env, ctx->cmd_exit_status));
int status;
int wpid = waitpid(ctx->pid, &status, WNOHANG);
if (wpid == ctx->pid) {
ctx->pid = CMD_EXIT;
ctx->cmd_exit_status = status;
return MAKE_OK(enif_make_int(env, status));
} else if (wpid != 0) {
perror("waitpid()");
}
ERL_NIF_TERM term = enif_make_tuple2(env, enif_make_int(env, wpid),
enif_make_int(env, status));
return MAKE_ERROR(term);
}
static int on_load(ErlNifEnv *env, void **priv, ERL_NIF_TERM load_info) {
struct ExilePriv *data = enif_alloc(sizeof(struct ExilePriv));
if (!data)
return 1;
data->rt =
enif_open_resource_type_x(env, "exile_resource", &rt_init,
ERL_NIF_RT_CREATE | ERL_NIF_RT_TAKEOVER, NULL);
ATOM_OK = enif_make_atom(env, "ok");
ATOM_ERROR = enif_make_atom(env, "error");
ATOM_UNDEFINED = enif_make_atom(env, "undefined");
ATOM_INVALID_CTX = enif_make_atom(env, "invalid_exile_exec_ctx");
ATOM_INVALID_CTX = enif_make_atom(env, "closed_pipe");
*priv = (void *)data;
return 0;
}
static void on_unload(ErlNifEnv *env, void *priv) {
debug("exile unload");
enif_free(priv);
}
static ErlNifFunc nif_funcs[] = {
{"exec_proc", 2, exec_proc, 0}, {"write_proc", 2, write_proc, 0},
{"read_proc", 2, read_proc, 0}, {"close_pipe", 2, close_pipe, 0},
{"terminate_proc", 1, terminate_proc, 0}, {"wait_proc", 1, wait_proc, 0},
{"kill_proc", 1, kill_proc, 0}, {"is_alive", 1, is_alive, 0},
};
-ERL_NIF_INIT(Elixir.Exile.ProcessHelper, nif_funcs, &on_load, NULL, NULL,
+ERL_NIF_INIT(Elixir.Exile.ProcessNif, nif_funcs, &on_load, NULL, NULL,
&on_unload)
diff --git a/lib/exile/process.ex b/lib/exile/process.ex
index ba77e2d..268a22c 100644
--- a/lib/exile/process.ex
+++ b/lib/exile/process.ex
@@ -1,331 +1,339 @@
defmodule Exile.Process do
- alias Exile.ProcessHelper
+ alias Exile.ProcessNif
require Logger
use GenServer
defmacro eagain(), do: 35
# delay between retries when io is busy (in milliseconds)
@default_opts %{io_busy_wait: 1, stderr_to_console: false}
def start_link(cmd, args, opts \\ %{}) do
opts = Map.merge(@default_opts, opts)
GenServer.start(__MODULE__, %{cmd: cmd, args: args, opts: opts})
end
def close_stdin(process) do
GenServer.call(process, :close_stdin, :infinity)
end
def write(process, binary) do
GenServer.call(process, {:write, binary}, :infinity)
end
def read(process, bytes) do
GenServer.call(process, {:read, bytes}, :infinity)
end
def kill(process, signal) when signal in [:sigkill, :sigterm] do
GenServer.call(process, {:kill, signal}, :infinity)
end
def await_exit(process, timeout \\ :infinity) do
GenServer.call(process, {:await_exit, timeout}, :infinity)
end
def stop(process), do: GenServer.call(process, :stop, :infinity)
## Server
defmodule Pending do
defstruct bin: [], remaining: 0, client_pid: nil
end
defstruct [
:cmd,
:cmd_args,
:opts,
:errno,
:context,
:status,
await: %{},
pending_read: nil,
pending_write: nil
]
alias __MODULE__
def init(%{cmd: cmd, args: args, opts: opts}) do
path = :os.find_executable(to_charlist(cmd))
unless path do
raise "Command not found: #{cmd}"
end
state = %__MODULE__{
cmd: path,
cmd_args: args,
opts: opts,
errno: nil,
status: :init,
await: %{},
pending_read: %Pending{},
pending_write: %Pending{}
}
{:ok, state, {:continue, nil}}
end
def handle_continue(nil, state) do
exec_args = Enum.map(state.cmd_args, &to_charlist/1)
stderr_to_console = if state.opts.stderr_to_console, do: 1, else: 0
- case ProcessHelper.exec_proc([state.cmd | exec_args], stderr_to_console) do
+ case ProcessNif.exec_proc([state.cmd | exec_args], stderr_to_console) do
{:ok, context} ->
start_watcher(context)
{:noreply, %Process{state | context: context, status: :start}}
{:error, errno} ->
raise "Failed to start command: #{state.cmd}, errno: #{errno}"
end
end
def handle_call(:stop, _from, state) do
- if ProcessHelper.is_alive(state.context) do
- do_close(state, :stdin)
- do_close(state, :stdout)
- do_kill(state.context, :sigkill)
- {:stop, :process_killed, :ok, %{state | status: {:exit, :killed}}}
- else
- {:stop, :normal, :ok, state}
- end
+ # watcher will terminate the external process
+ {:stop, :normal, :ok, state}
end
def handle_call(_, _from, %{status: {:exit, status}}), do: {:reply, {:error, {:exit, status}}}
def handle_call({:await_exit, timeout}, from, state) do
tref =
if timeout != :infinity do
Elixir.Process.send_after(self(), {:await_exit_timeout, from}, timeout)
else
nil
end
state = put_timer(state, from, :timeout, tref)
check_exit(state, from)
end
def handle_call({:write, binary}, from, state) do
pending = %Pending{bin: binary, client_pid: from}
do_write(%Process{state | pending_write: pending})
end
def handle_call({:read, bytes}, from, state) do
pending = %Pending{remaining: bytes, client_pid: from}
do_read(%Process{state | pending_read: pending})
end
def handle_call(:close_stdin, _from, state), do: do_close(state, :stdin)
def handle_call({:kill, signal}, _from, state) do
do_kill(state.context, signal)
{:reply, :ok, %{state | status: {:exit, :killed}}}
end
def handle_info({:check_exit, from}, state), do: check_exit(state, from)
def handle_info({:await_exit_timeout, from}, state) do
cancel_timer(state, from, :check)
receive do
{:check_exit, ^from} -> :ok
after
0 -> :ok
end
GenServer.reply(from, :timeout)
{:noreply, clear_await(state, from)}
end
def handle_info({:select, context, _ref, :ready_output}, state) do
do_write(%Process{state | context: context})
end
def handle_info({:select, context, _ref, :ready_input}, state) do
do_read(%Process{state | context: context})
end
def handle_info(msg, _state), do: raise(msg)
defp do_write(%Process{pending_write: pending} = state) do
- case ProcessHelper.write_proc(state.context, pending.bin) do
+ case ProcessNif.write_proc(state.context, pending.bin) do
{:ok, size} ->
if size < IO.iodata_length(pending.bin) do
binary = IO.iodata_to_binary(pending.bin)
binary = binary_part(binary, size, IO.iodata_length(pending.bin) - size)
{:noreply, %{state | pending_write: %Pending{bin: binary}}}
else
GenServer.reply(pending.client_pid, :ok)
{:noreply, %{state | pending_write: %Pending{}}}
end
{:error, eagain()} ->
{:noreply, state}
{:error, errno} ->
GenServer.reply(pending.client_pid, {:error, errno})
{:noreply, %{state | errno: errno}}
end
end
defp do_read(%Process{pending_read: %Pending{remaining: nil} = pending} = state) do
- case ProcessHelper.read_proc(state.context, -1) do
+ case ProcessNif.read_proc(state.context, -1) do
{:ok, <<>>} ->
GenServer.reply(pending.client_pid, {:eof, []})
{:noreply, state}
{:ok, binary} ->
GenServer.reply(pending.client_pid, {:ok, binary})
{:noreply, state}
{:error, eagain()} ->
{:noreply, state}
{:error, errno} ->
GenServer.reply(pending.client_pid, {:error, errno})
{:noreply, %{state | errno: errno}}
end
end
defp do_read(%Process{pending_read: pending} = state) do
- case ProcessHelper.read_proc(state.context, pending.remaining) do
+ case ProcessNif.read_proc(state.context, pending.remaining) do
{:ok, <<>>} ->
GenServer.reply(pending.client_pid, {:eof, pending.bin})
{:noreply, %Process{state | pending_read: %Pending{}}}
{:ok, binary} ->
if IO.iodata_length(binary) < pending.remaining do
pending = %Pending{
pending
| bin: [pending.bin | binary],
remaining: pending.remaining - IO.iodata_length(binary)
}
{:noreply, %Process{state | pending_read: pending}}
else
GenServer.reply(pending.client_pid, {:ok, [state.pending_read.bin | binary]})
{:noreply, %Process{state | pending_read: %Pending{}}}
end
{:error, eagain()} ->
{:noreply, state}
{:error, errno} ->
GenServer.reply(pending.client_pid, {:error, errno})
{:noreply, %{state | pending_read: %Pending{}, errno: errno}}
end
end
defp check_exit(state, from) do
- case ProcessHelper.wait_proc(state.context) do
+ case ProcessNif.wait_proc(state.context) do
{:ok, status} ->
GenServer.reply(from, {:ok, status})
cancel_timer(state, from, :timeout)
{:noreply, clear_await(state, from)}
{:error, {0, _}} ->
# Ideally we should not poll and we should handle this with SIGCHLD signal
tref = Elixir.Process.send_after(self(), {:check_exit, from}, state.opts.io_busy_wait)
{:noreply, put_timer(state, from, :check, tref)}
{:error, {-1, status}} ->
GenServer.reply(from, {:error, status})
cancel_timer(state, from, :timeout)
{:noreply, clear_await(state, from)}
end
end
- defp do_kill(context, :sigkill), do: ProcessHelper.kill_proc(context)
+ defp do_kill(context, :sigkill), do: ProcessNif.kill_proc(context)
- defp do_kill(context, :sigterm), do: ProcessHelper.terminate_proc(context)
+ defp do_kill(context, :sigterm), do: ProcessNif.terminate_proc(context)
defp do_close(state, type) do
- case ProcessHelper.close_pipe(state.context, stream_type(type)) do
+ case ProcessNif.close_pipe(state.context, stream_type(type)) do
:ok ->
{:reply, :ok, state}
{:error, errno} ->
raise errno
{:reply, {:error, errno}, %Process{state | errno: errno}}
end
end
defp clear_await(state, from) do
%Process{state | await: Map.delete(state.await, from)}
end
defp cancel_timer(state, from, key) do
case get_timer(state, from, key) do
nil -> :ok
tref -> Elixir.Process.cancel_timer(tref)
end
end
defp put_timer(state, from, key, timer) do
if Map.has_key?(state.await, from) do
await = put_in(state.await, [from, key], timer)
%Process{state | await: await}
else
%Process{state | await: %{from => %{key => timer}}}
end
end
defp get_timer(state, from, key), do: get_in(state.await, [from, key])
- @stdin_close_wait 3000
- @sigterm_wait 1000
-
# Try to gracefully terminate external proccess if the genserver associated with the process is killed
defp start_watcher(context) do
process_server = self()
watcher_pid = spawn(fn -> watcher(process_server, context) end)
receive do
{^watcher_pid, :done} -> :ok
end
end
defp stream_type(:stdin), do: 0
defp stream_type(:stdout), do: 1
+ defp process_exit?(context) do
+ match?({:ok, _}, ProcessNif.wait_proc(context))
+ end
+
+ defp process_exit?(context, timeout) do
+ if process_exit?(context) do
+ true
+ else
+ :timer.sleep(timeout)
+ process_exit?(context)
+ end
+ end
+
+ # for proper process exit parent of the child *must* wait() for
+ # child processes termination exit and "pickup" after the exit
+ # (receive child exit_status). Resources acquired by child such as
+ # file descriptors won't be released even if the child process
+ # itself is terminated.
defp watcher(process_server, context) do
ref = Elixir.Process.monitor(process_server)
send(process_server, {self(), :done})
receive do
- {:DOWN, ^ref, :process, ^process_server, :normal} ->
- :ok
-
{:DOWN, ^ref, :process, ^process_server, _reason} ->
- case ProcessHelper.wait_proc(context) do
- {:ok, _status} ->
- # TODO: check stauts
- nil
-
- {:error, {_, _}} ->
- Logger.debug(fn -> "Killing" end)
-
- with _ <- ProcessHelper.close_pipe(context, stream_type(:stdin)),
- _ <- ProcessHelper.close_pipe(context, stream_type(:stdout)),
- _ <- :timer.sleep(@stdin_close_wait),
- {:error, _} <- ProcessHelper.wait_proc(context),
- _ <- ProcessHelper.terminate_proc(context),
- _ <- :timer.sleep(@sigterm_wait),
- {:error, _} <- ProcessHelper.wait_proc(context),
- _ <- ProcessHelper.kill_proc(context) do
- Logger.debug(fn -> "Killed process" end)
- end
+ try do
+ process_exit?(context) && throw(:done)
+
+ Logger.debug(fn -> "Killing" end)
+ ProcessNif.close_pipe(context, stream_type(:stdin))
+ ProcessNif.close_pipe(context, stream_type(:stdout))
+
+ # at max we wait for 100ms for program to exit
+ process_exit?(context, 100) && throw(:done)
+
+ ProcessNif.terminate_proc(context)
+ process_exit?(context, 100) && throw(:done)
+
+ ProcessNif.kill_proc(context)
+ process_exit?(context, 1000) && throw(:done)
+
+ Logger.error("[exile] failed to kill external process")
+ raise "Failed to kill external process"
+ catch
+ :done -> Logger.debug(fn -> "Exited external program successfully" end)
end
end
end
end
diff --git a/lib/exile/process_helper.ex b/lib/exile/process_nif.ex
similarity index 96%
rename from lib/exile/process_helper.ex
rename to lib/exile/process_nif.ex
index f8353de..af6a1de 100644
--- a/lib/exile/process_helper.ex
+++ b/lib/exile/process_nif.ex
@@ -1,24 +1,24 @@
-defmodule Exile.ProcessHelper do
+defmodule Exile.ProcessNif do
@on_load :load_nifs
def load_nifs do
nif_path = :filename.join(:code.priv_dir(:exile), "exile")
:erlang.load_nif(nif_path, 0)
end
def exec_proc(_cmd, _stderr_to_console), do: :erlang.nif_error(:nif_library_not_loaded)
def write_proc(_context, _bin), do: :erlang.nif_error(:nif_library_not_loaded)
def read_proc(_context, _bytes), do: :erlang.nif_error(:nif_library_not_loaded)
def close_pipe(_context, _pipe), do: :erlang.nif_error(:nif_library_not_loaded)
def kill_proc(_context), do: :erlang.nif_error(:nif_library_not_loaded)
def terminate_proc(_context), do: :erlang.nif_error(:nif_library_not_loaded)
def wait_proc(_context), do: :erlang.nif_error(:nif_library_not_loaded)
def is_alive(_context), do: :erlang.nif_error(:nif_library_not_loaded)
end
File Metadata
Details
Attached
Mime Type
text/x-diff
Expires
Sat, Aug 8, 5:59 PM (1 d, 17 h)
Storage Engine
blob
Storage Format
Raw Data
Storage Handle
1722679
Default Alt Text
(28 KB)
Attached To
Mode
R14 exile
Attached
Detach File
Event Timeline
Log In to Comment