Page MenuHomePhorge

No OneTemporary

Size
30 KB
Referenced Files
None
Subscribers
None
diff --git a/c_src/exile.c b/c_src/exile.c
index 68284b1..98d94fc 100644
--- a/c_src/exile.c
+++ b/c_src/exile.c
@@ -1,431 +1,506 @@
#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 ERL_OK(term) enif_make_tuple2(env, enif_make_atom(env, "ok"), term)
-#define ERL_ERROR(term) \
- enif_make_tuple2(env, enif_make_atom(env, "error"), term)
-static const int PIPE_READ = 0;
-static const int PIPE_WRITE = 1;
-static const int MAX_ARGUMENTS = 20;
+#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 ExecResults {
+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;
- pid_t pid;
- int pipe_in;
- int pipe_out;
+ ExecContext context;
} ExecResult;
-struct ExilePriv {
- /* ERL_NIF_TERM atom_ok; */
- /* ERL_NIF_TERM atom_undefined; */
-
- ErlNifResourceType *rt;
- /* void *read_resource; */
- /* void *write_resource; */
-};
-
static void rt_dtor(ErlNifEnv *env, void *obj) {
printf("----- rt_dtor called\n");
}
static void rt_stop(ErlNifEnv *env, void *obj, int fd, int is_direct_call) {
printf("----- rt_stop called\n");
}
static void rt_down(ErlNifEnv *env, void *obj, ErlNifPid *pid,
ErlNifMonitor *monitor) {
printf("----- rt_down called\n");
}
static ErlNifResourceTypeInit rt_init = {rt_dtor, rt_stop, rt_down};
-typedef struct ExecContext {
- int cmd_input_fd;
- int cmd_output_fd;
- pid_t pid;
-} ExecContext;
-
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(error) \
+#define RETURN_ERROR(__err) \
do { \
- fprintf(stderr, "error in start_proccess(), %s:%d %s\n", __FILE__, \
- __LINE__, strerror(errno)); \
- result.status = error; \
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)
}
- if (set_flag(pipes[STDIN_FILENO][PIPE_READ], O_CLOEXEC) < 0 ||
- set_flag(pipes[STDOUT_FILENO][PIPE_WRITE], O_CLOEXEC) < 0 ||
- set_flag(pipes[STDIN_FILENO][PIPE_WRITE], O_CLOEXEC | O_NONBLOCK) < 0 ||
- set_flag(pipes[STDOUT_FILENO][PIPE_READ], O_CLOEXEC | O_NONBLOCK) < 0) {
+ 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:
+ case 0: // child
+
+ // close default stdio fd
close(STDIN_FILENO);
close(STDOUT_FILENO);
- if (dup2(pipes[STDIN_FILENO][PIPE_READ], STDIN_FILENO) < 0)
+ if (dup2(r_cmdin, STDIN_FILENO) < 0) {
+ debug("failed dup command input pipe to stdin");
RETURN_ERROR(PIPE_DUP_ERROR)
- if (dup2(pipes[STDOUT_FILENO][PIPE_WRITE], STDOUT_FILENO) < 0)
+ }
+ 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)
+ 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:
- close(pipes[STDIN_FILENO][PIPE_READ]);
- close(pipes[STDOUT_FILENO][PIPE_WRITE]);
- result.pid = pid;
- result.pipe_in = pipes[STDIN_FILENO][PIPE_WRITE];
- result.pipe_out = pipes[STDOUT_FILENO][PIPE_READ];
+ 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);
- ERL_NIF_TERM ret;
ExecContext *ctx = NULL;
switch (result.status) {
case SUCCESS:
ctx = enif_alloc_resource(data->rt, sizeof(ExecContext));
- ctx->cmd_input_fd = result.pipe_in;
- ctx->cmd_output_fd = result.pipe_out;
- ctx->pid = result.pid;
+ ctx->cmd_input_fd = result.context.cmd_input_fd;
+ ctx->cmd_output_fd = result.context.cmd_output_fd;
+ ctx->pid = result.context.pid;
- printf("cmd_in: %d cmd_out: %d pid: %d\n", result.pipe_in, result.pipe_out,
- result.pid);
+ debug("pid: %d cmd_in_fd: %d cmd_out_fd: %d", ctx->pid, ctx->cmd_input_fd,
+ ctx->cmd_output_fd);
// TODO: exit the command gracefully when resource is released by GC
/* enif_release_resource(ctx); */
- return ERL_OK(enif_make_resource(env, ctx));
+ return MAKE_OK(enif_make_resource(env, ctx));
default:
- ret = enif_make_int(env, result.err);
- return ERL_ERROR(ret);
+ 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);
- struct ExilePriv *data = enif_priv_data(env);
ExecContext *ctx = NULL;
- if (enif_get_resource(env, argv[0], data->rt, (void **)&ctx) == false) {
- return enif_make_badarg(env);
- }
+ 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 ERL_OK(enif_make_int(env, result));
+ 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 ERL_ERROR(enif_make_int(env, retval));
- return ERL_OK(enif_make_int(env, result));
+ 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 ERL_ERROR(enif_make_int(env, retval));
- return ERL_ERROR(enif_make_int(env, EAGAIN));
+ return MAKE_ERROR(enif_make_int(env, retval));
+ return MAKE_ERROR(enif_make_int(env, EAGAIN));
} else { // Error
perror("write()");
- return ERL_ERROR(enif_make_int(env, errno));
+ return MAKE_ERROR(enif_make_int(env, errno));
}
}
static ERL_NIF_TERM close_pipe(ErlNifEnv *env, int argc,
const ERL_NIF_TERM argv[]) {
- struct ExilePriv *data = enif_priv_data(env);
ExecContext *ctx = NULL;
- if (enif_get_resource(env, argv[0], data->rt, (void **)&ctx) == false) {
- return enif_make_badarg(env);
- }
+ GET_CTX(env, argv[0], ctx);
int kind;
enif_get_int(env, argv[1], &kind);
int result;
switch (kind) {
case 0:
- result = close(ctx->cmd_input_fd);
- break;
+ 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:
- result = close(ctx->cmd_output_fd);
- break;
+ 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);
}
-
- if (result == 0) {
- return enif_make_atom(env, "ok");
- } else {
- perror("close()");
- return ERL_ERROR(enif_make_int(env, errno));
- }
}
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);
- struct ExilePriv *data = enif_priv_data(env);
ExecContext *ctx = NULL;
- if (enif_get_resource(env, argv[0], data->rt, (void **)&ctx) == false) {
- return enif_make_badarg(env);
- }
+ 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 ERL_OK(bin_term);
+ return MAKE_OK(bin_term);
} else if (result > 0) { // request partially satisfied
int retval = select_read(env, ctx);
if (retval != 0)
- return ERL_ERROR(enif_make_int(env, retval));
- return ERL_OK(bin_term);
+ return MAKE_ERROR(enif_make_int(env, retval));
+ return MAKE_OK(bin_term);
} else if (result == 0) { // EOF
- return ERL_OK(bin_term);
+ return MAKE_OK(bin_term);
} else if (errno == EAGAIN) { // busy
int retval = select_read(env, ctx);
if (retval != 0)
- return ERL_ERROR(enif_make_int(env, retval));
- return ERL_ERROR(enif_make_int(env, EAGAIN));
+ return MAKE_ERROR(enif_make_int(env, retval));
+ return MAKE_ERROR(enif_make_int(env, EAGAIN));
} else { // Error
perror("read()");
- return ERL_ERROR(enif_make_int(env, errno));
+ return MAKE_ERROR(enif_make_int(env, errno));
}
}
static ERL_NIF_TERM is_alive(ErlNifEnv *env, int argc,
const ERL_NIF_TERM argv[]) {
- struct ExilePriv *data = enif_priv_data(env);
ExecContext *ctx = NULL;
- if (enif_get_resource(env, argv[0], data->rt, (void **)&ctx) == false) {
- return enif_make_badarg(env);
- }
+ 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[]) {
- struct ExilePriv *data = enif_priv_data(env);
ExecContext *ctx = NULL;
- if (enif_get_resource(env, argv[0], data->rt, (void **)&ctx) == false) {
- return enif_make_badarg(env);
- }
+ 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[]) {
- struct ExilePriv *data = enif_priv_data(env);
ExecContext *ctx = NULL;
- if (enif_get_resource(env, argv[0], data->rt, (void **)&ctx) == false) {
- return enif_make_badarg(env);
- }
+ 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[]) {
- struct ExilePriv *data = enif_priv_data(env);
ExecContext *ctx = NULL;
- if (enif_get_resource(env, argv[0], data->rt, (void **)&ctx) == false) {
- return enif_make_badarg(env);
- }
+ 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) {
- return ERL_OK(enif_make_int(env, status));
- } else {
+ 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 ERL_ERROR(term);
}
+ ERL_NIF_TERM term = enif_make_tuple2(env, enif_make_int(env, wpid),
+ enif_make_int(env, status));
+ return MAKE_ERROR(term);
+
}
-static int load(ErlNifEnv *env, void **priv, ERL_NIF_TERM load_info) {
+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->atom_ok = enif_make_atom(env, "ok"); */
- /* data->atom_undefined = enif_make_atom(env, "undefined"); */
+ data->rt =
+ enif_open_resource_type_x(env, "exile_resource", &rt_init,
+ ERL_NIF_RT_CREATE | ERL_NIF_RT_TAKEOVER, NULL);
- data->rt = enif_open_resource_type_x(env, "exile_resource", &rt_init,
- ERL_NIF_RT_CREATE, 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, &load, NULL, NULL, NULL)
+ERL_NIF_INIT(Elixir.Exile.ProcessHelper, nif_funcs, &on_load, NULL, NULL,
+ &on_unload)
diff --git a/lib/exile/process.ex b/lib/exile/process.ex
index 705f230..ba77e2d 100644
--- a/lib/exile/process.ex
+++ b/lib/exile/process.ex
@@ -1,340 +1,331 @@
defmodule Exile.Process do
alias Exile.ProcessHelper
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,
- :stdin_closed,
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
{: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
- # do_close(state, :stdin)
- # do_close(state, :stdout)
-
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
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, %Process{stdin_closed: true} = state),
- do: {:reply, {:error, :closed}, state}
-
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
{: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
{: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
{: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
{: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, :sigterm), do: ProcessHelper.terminate_proc(context)
- defp do_close(%Process{stdin_closed: true} = state, type) do
- {:reply, :ok, state}
- end
-
defp do_close(state, type) do
case ProcessHelper.close_pipe(state.context, stream_type(type)) do
:ok ->
- {:reply, :ok, %Process{state | stdin_closed: true}}
+ {: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 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
end
end
end
end

File Metadata

Mime Type
text/x-diff
Expires
Sat, Aug 8, 7:01 PM (1 d, 13 h)
Storage Engine
blob
Storage Format
Raw Data
Storage Handle
1722694
Default Alt Text
(30 KB)

Event Timeline