diff --git a/CHANGELOG.md b/CHANGELOG.md index 1514ebf..fd19b18 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,42 @@ # Changelog +## Unreleased + +### Fixed + +- Tasks on the event loop pool that awaited `asyncio.sleep` for more than a + few milliseconds never completed when several ran at once: with 24 + concurrent 50 ms sleeps, 23 timed out, while the same tasks run one after + another always finished (in Hornbeam, concurrent ASGI requests hung until + the request timeout). Every pool loop scheduled its timers on the default + loop and polled that loop's queue, where the per-loop callback ids + collided. Each pool loop now drives a Python loop bound to its own + resource, a timer expiry is dispatched to the loop that set it, and a loop + with no tasks leaves its pending events to whoever polls it. +- A Python function called with `py:call` that called `erlang.call`, where + the Erlang callback called `py:call` again, hung until the request timeout + and then failed with "callback synchronisation lost; retry" or + `{error, timeout}`, even one level deep; the same chain started with + `py:eval` worked. `erlang.call` blocked the context thread on the thread + worker pipe and the nested `py:call` waited on that context. The context + thread now serves its own requests while it waits for the callback, so the + nested `py:call` runs inline on the same context, at any depth and with + several sequential `erlang.call`s in one function; the Python code still + runs exactly once. Inside a running asyncio loop (`py_context:start_loop/2`, + `asyncio.run` in a called function) the thread path stays. +- The shared buffer's flow control and the sync `erlang.sleep` call Erlang + through a blocking path that never suspends the request, so a `py:eval` + that reads a shared buffer or sleeps is not replayed around them. +- After a callback resumed, a `py:eval` or `py:call` that used a function + defined with `py:exec` failed with `NameError: name ... is not defined`: + the replay ran in the context globals instead of the caller's namespace. + The replay now runs in the namespace of the original request. +- An Erlang callback could not see the caller's `__main__` functions and + was routed to another context: `py:call('__main__', double, [X])` from a + callback failed with "module '__main__' has no attribute 'double'". The + callback process is now bound to the suspended context and shares the + caller's namespace, so the README reentrant example passes as written. + ## 5.0.0 (2026-08-29) ### Added diff --git a/README.md b/README.md index ca06e00..d82f198 100644 --- a/README.md +++ b/README.md @@ -197,9 +197,10 @@ def process(x): %% 10 → double_via_python → double(10)=20 → +1 = 21 ``` -The implementation uses a suspension/resume mechanism that frees the dirty -scheduler while the Erlang callback executes, preventing deadlocks even with -multiple levels of nesting. +While the Erlang callback runs, the context thread keeps serving requests +for its context, so the callback's own `py:call` runs there and nesting +works at any depth. The callback shares the caller's Python namespace, +which is how it can see `double`. No dirty scheduler is held meanwhile. ## Shared State Between Workers diff --git a/c_src/py_callback.c b/c_src/py_callback.c index 772d817..6ac4126 100644 --- a/c_src/py_callback.c +++ b/c_src/py_callback.c @@ -365,7 +365,8 @@ static suspended_context_state_t *create_suspended_context_state_for_call( ErlNifBinary *module_bin, ErlNifBinary *func_bin, ERL_NIF_TERM args_term, - ERL_NIF_TERM kwargs_term) { + ERL_NIF_TERM kwargs_term, + py_env_resource_t *penv) { /* Allocate the suspended context state resource */ suspended_context_state_t *state = enif_alloc_resource( @@ -379,6 +380,10 @@ static suspended_context_state_t *create_suspended_context_state_for_call( state->ctx = ctx; enif_keep_resource(ctx); /* Keep ctx alive while suspended state exists */ + state->penv = penv; + if (penv != NULL) { + enif_keep_resource(penv); /* The replay runs in the request's namespace */ + } state->callback_id = tl_pending_callback_id; state->request_type = PY_REQ_CALL; @@ -450,7 +455,8 @@ static suspended_context_state_t *create_suspended_context_state_for_eval( ErlNifEnv *env, py_context_t *ctx, ErlNifBinary *code_bin, - ERL_NIF_TERM locals_term) { + ERL_NIF_TERM locals_term, + py_env_resource_t *penv) { (void)env; @@ -466,6 +472,10 @@ static suspended_context_state_t *create_suspended_context_state_for_eval( state->ctx = ctx; enif_keep_resource(ctx); /* Keep ctx alive while suspended state exists */ + state->penv = penv; + if (penv != NULL) { + enif_keep_resource(penv); /* The replay runs in the request's namespace */ + } state->callback_id = tl_pending_callback_id; state->request_type = PY_REQ_EVAL; @@ -1541,6 +1551,79 @@ static PyObject *py_consume_time_slice(PyObject *self, PyObject *args) { * * This allows the dirty scheduler to be freed while waiting for the callback. */ +/** + * True when an asyncio loop is running on this thread. + * + * A suspension unwinds the whole Python stack of the request, which cannot + * be done from inside a running loop: erlang._run_loop_forever (the + * context worker loop) and asyncio.run() inside a called function keep the + * blocking thread path for their erlang.call. + */ +static bool asyncio_loop_running_here(void) { + PyObject *modules = PyImport_GetModuleDict(); + if (modules == NULL) { + return false; + } + PyObject *events = PyDict_GetItemString(modules, "asyncio.events"); /* Borrowed */ + if (events == NULL) { + return false; /* asyncio never imported: no loop can be running */ + } + PyObject *get_running = PyObject_GetAttrString(events, "_get_running_loop"); + if (get_running == NULL) { + PyErr_Clear(); + return false; + } + PyObject *loop = PyObject_CallNoArgs(get_running); + Py_DECREF(get_running); + if (loop == NULL) { + PyErr_Clear(); + return false; + } + bool running = (loop != Py_None); + Py_DECREF(loop); + return running; +} + +/** + * erlang._call_blocking(name, *args): the thread worker path, always. + * + * Blocks the calling Python thread on the thread worker pipe until the + * Erlang function returns; the Python frame is never suspended or + * replayed. For the library's own flow-control callbacks (shared buffer + * waits, erlang.sleep), whose callers keep state across the call, and for + * any Python thread. + */ +static PyObject *erlang_call_blocking_impl(PyObject *self, PyObject *args) { + (void)self; + + Py_ssize_t nargs = PyTuple_Size(args); + if (nargs < 1) { + PyErr_SetString(PyExc_TypeError, "erlang.call requires at least a function name"); + return NULL; + } + + PyObject *name_obj = PyTuple_GetItem(args, 0); + if (!PyUnicode_Check(name_obj)) { + PyErr_SetString(PyExc_TypeError, "Function name must be a string"); + return NULL; + } + const char *func_name = PyUnicode_AsUTF8(name_obj); + if (func_name == NULL) { + return NULL; + } + size_t func_name_len = strlen(func_name); + + /* Build args list (remaining args) */ + PyObject *call_args = PyTuple_GetSlice(args, 1, nargs); + if (call_args == NULL) { + return NULL; + } + + PyObject *result = thread_worker_call(func_name, func_name_len, call_args); + Py_DECREF(call_args); + return result; +} + static PyObject *erlang_call_impl(PyObject *self, PyObject *args) { (void)self; @@ -1560,58 +1643,62 @@ static PyObject *erlang_call_impl(PyObject *self, PyObject *args) { /* * Check if we have a callback handler available. * Priority: - * 1. tl_current_context with suspension enabled (new process-per-context API) - * 2. tl_current_context with callback_handler (old blocking pipe mode) - * 3. thread_worker_call (spawned threads) + * 1. tl_current_context with suspension enabled: eval requests and + * resume replays (ctx_execute_eval_with_env, nif_context_resume) + * 2. tl_current_context with a request in flight: call requests wait + * inline, serving the context's queue (ctx_call_erlang_inline) + * 3. tl_current_context with callback_handler (old blocking pipe mode) + * 4. thread_worker_call (spawned threads) * - * NOTE: In OWN_GIL mode, erlang.call() goes through thread_worker_call() - * rather than using suspension/resume. This is because OWN_GIL contexts - * bypass the suspension protocol - the dedicated pthread that owns the GIL - * cannot be suspended. As a result, the call executes on a different - * context/interpreter (the thread worker), not the calling OWN_GIL context. - * Re-entrant calls back to the same OWN_GIL context are not supported. + * Inside a running asyncio loop the request stack cannot be unwound or + * blocked without stalling the loop: those calls take path 4. */ - bool has_context_suspension = (tl_current_context != NULL && tl_allow_suspension); + bool loop_running = asyncio_loop_running_here(); + bool has_context_suspension = (tl_current_context != NULL && tl_allow_suspension && + !loop_running); bool has_context_handler = (tl_current_context != NULL && tl_current_context->has_callback_handler); + bool has_context_inline = (tl_current_context != NULL && !tl_allow_suspension && + tl_current_context->has_current_caller && + !has_context_handler && !loop_running); - if (!has_context_suspension && !has_context_handler) { - /* - * Not an executor thread - use thread worker path. - * This enables any spawned Python thread to call erlang.call(): - * - threading.Thread instances - * - concurrent.futures.ThreadPoolExecutor workers - * - Any other Python threads - * - OWN_GIL contexts (which don't support suspension) - */ + if (has_context_inline) { Py_ssize_t nargs = PyTuple_Size(args); if (nargs < 1) { PyErr_SetString(PyExc_TypeError, "erlang.call requires at least a function name"); return NULL; } - PyObject *name_obj = PyTuple_GetItem(args, 0); if (!PyUnicode_Check(name_obj)) { PyErr_SetString(PyExc_TypeError, "Function name must be a string"); return NULL; } - const char *func_name = PyUnicode_AsUTF8(name_obj); + Py_ssize_t func_name_len = 0; + const char *func_name = PyUnicode_AsUTF8AndSize(name_obj, &func_name_len); if (func_name == NULL) { return NULL; } - size_t func_name_len = strlen(func_name); - - /* Build args list (remaining args) */ PyObject *call_args = PyTuple_GetSlice(args, 1, nargs); if (call_args == NULL) { return NULL; } - - /* Use thread worker call */ - PyObject *result = thread_worker_call(func_name, func_name_len, call_args); + PyObject *result = ctx_call_erlang_inline(tl_current_context, func_name, + (size_t)func_name_len, call_args); Py_DECREF(call_args); return result; } + if (!has_context_suspension && !has_context_handler) { + /* + * Not an executor thread - use thread worker path. + * This enables any spawned Python thread to call erlang.call(): + * - threading.Thread instances + * - concurrent.futures.ThreadPoolExecutor workers + * - Any other Python threads + * - code running inside an asyncio loop on a context thread + */ + return erlang_call_blocking_impl(self, args); + } + Py_ssize_t nargs = PyTuple_Size(args); if (nargs < 1) { PyErr_SetString(PyExc_TypeError, "erlang.call requires at least a function name"); @@ -3189,6 +3276,8 @@ static PyObject *erlang_whereis_impl(PyObject *self, PyObject *args) { /* Python method definitions for erlang module */ static PyMethodDef ErlangModuleMethods[] = { + {"_call_blocking", erlang_call_blocking_impl, METH_VARARGS, + "Call an Erlang function on the blocking thread path (no suspension, no replay)"}, {"call", erlang_call_impl, METH_VARARGS, "Call a registered Erlang function.\n\n" "Usage: erlang.call('func_name', arg1, arg2, ...)\n" diff --git a/c_src/py_event_loop.c b/c_src/py_event_loop.c index f88d46f..fe4080a 100644 --- a/c_src/py_event_loop.c +++ b/c_src/py_event_loop.c @@ -95,6 +95,15 @@ ERL_NIF_TERM ATOM_DISPATCH; /** @brief Name for the PyCapsule storing event loop pointer */ static const char *EVENT_LOOP_CAPSULE_NAME = "erlang_python.event_loop"; +/* Capsule over a loop resource for ErlangEventLoop (defined with the + * capsule helpers below). */ +static PyObject *make_loop_capsule(erlang_event_loop_t *loop); + +/* Timer refs are keyed in the receiving worker's map. Loops without a + * worker of their own share the global worker, so the ids must be unique + * across loops, not per loop. */ +static _Atomic uint64_t g_next_timer_ref = 1; + /** @brief Module attribute name for storing the event loop */ static const char *EVENT_LOOP_ATTR_NAME = "_loop"; @@ -2862,6 +2871,41 @@ static void loop_gil_release(erlang_event_loop_t *loop, loop_gil_t *g) { PyGILState_Release(g->gstate); } +/** + * event_loop_release_python_loop(LoopRef) -> ok | {error, Reason} + * + * Drops the loop's reference to its Python ErlangEventLoop. That object + * holds a capsule that keeps the loop resource alive, so without this the + * two keep each other forever once the loop was driven by + * process_ready_tasks. Called by the pool before event_loop_destroy. + * Dirty: takes the loop's GIL. + */ +ERL_NIF_TERM nif_event_loop_release_python_loop(ErlNifEnv *env, int argc, + const ERL_NIF_TERM argv[]) { + (void)argc; + + erlang_event_loop_t *loop; + if (!enif_get_resource(env, argv[0], EVENT_LOOP_RESOURCE_TYPE, + (void **)&loop)) { + return make_error(env, "invalid_loop"); + } + if (loop->py_loop == NULL) { + return ATOM_OK; + } + if (!runtime_is_running()) { + return make_error(env, "python_not_running"); + } + + loop_gil_t gil; + if (!loop_gil_acquire(loop, &gil)) { + return make_error(env, "interpreter_gone"); + } + Py_CLEAR(loop->py_loop); + loop->py_loop_valid = false; + loop_gil_release(loop, &gil); + return ATOM_OK; +} + void event_loop_detach_interpreter(erlang_event_loop_t *loop) { if (loop == NULL) { return; @@ -3118,6 +3162,13 @@ ERL_NIF_TERM nif_process_ready_tasks(ErlNifEnv *env, int argc, /* Lazy loop creation (uvloop-style): create Python loop on first use */ if (!loop->py_loop_valid || loop->py_loop == NULL) { + if (num_tasks == 0) { + /* Nothing to schedule and no Python loop to run callbacks in: + * the pending events stay queued for whoever polls this loop + * (get_pending, or a Python loop attached later). */ + loop_gil_release(loop, &gil); + return ATOM_OK; + } /* Create ErlangEventLoop directly instead of via asyncio.new_event_loop(). * This is necessary because dirty NIF scheduler threads don't have the * event loop policy set. asyncio.new_event_loop() would create a @@ -3147,7 +3198,22 @@ ERL_NIF_TERM nif_process_ready_tasks(ErlNifEnv *env, int argc, return make_error(env, "loop_class_not_found"); } - PyObject *new_loop = PyObject_CallNoArgs(loop_class); + /* The interpreter's default loop keeps the global capsule path in + * ErlangEventLoop.__init__. Any other loop (the event loop pool) + * gets a capsule over its own resource, so its timers, fd events + * and pending queue stay on this loop instead of the global one. */ + PyObject *new_loop; + if (loop == get_interpreter_event_loop()) { + new_loop = PyObject_CallNoArgs(loop_class); + } else { + PyObject *capsule = make_loop_capsule(loop); + if (capsule == NULL) { + new_loop = NULL; + } else { + new_loop = PyObject_CallOneArg(loop_class, capsule); + Py_DECREF(capsule); + } + } Py_DECREF(loop_class); if (new_loop == NULL) { PyErr_Clear(); @@ -6963,7 +7029,7 @@ static PyObject *py_schedule_timer(PyObject *self, PyObject *args) { } if (delay_ms < 0) delay_ms = 0; - uint64_t timer_ref_id = atomic_fetch_add(&loop->next_callback_id, 1); + uint64_t timer_ref_id = atomic_fetch_add(&g_next_timer_ref, 1); /* Use per-call env for thread safety in free-threaded Python */ ErlNifEnv *msg_env = enif_alloc_env(); @@ -7216,6 +7282,21 @@ static void global_loop_capsule_destructor(PyObject *capsule) { } } +/** + * Capsule over an Erlang-owned loop resource. The capsule keeps the + * resource alive and releases it when collected; it never signals + * shutdown, Erlang does that through event_loop_destroy. + */ +static PyObject *make_loop_capsule(erlang_event_loop_t *loop) { + enif_keep_resource(loop); + PyObject *capsule = PyCapsule_New(loop, LOOP_CAPSULE_NAME, + global_loop_capsule_destructor); + if (capsule == NULL) { + enif_release_resource(loop); + } + return capsule; +} + /* Python function: _loop_new() -> capsule */ static PyObject *py_loop_new(PyObject *self, PyObject *args) { (void)self; @@ -7470,10 +7551,7 @@ static PyObject *py_get_global_loop_capsule(PyObject *self, PyObject *args) { return NULL; } - /* Keep the resource alive while capsule exists */ - enif_keep_resource(loop); - - return PyCapsule_New(loop, LOOP_CAPSULE_NAME, global_loop_capsule_destructor); + return make_loop_capsule(loop); } /** @@ -8044,27 +8122,18 @@ static PyObject *py_schedule_timer_for(PyObject *self, PyObject *args) { return NULL; } - /* For timer scheduling, we need to use the global interpreter loop which - * has the worker process. The capsule's loop may be a Python-created loop - * that doesn't have has_worker set, which would cause timer dispatches - * to go to the router instead of the worker, breaking the event loop flow. - * - * The global loop (created by Erlang) has has_worker=true and its worker - * properly triggers process_ready_tasks after timer dispatch. */ - erlang_event_loop_t *target_loop = get_interpreter_event_loop(); - if (target_loop == NULL) { - /* Fall back to capsule's loop if global not available */ - target_loop = loop; - } - - if (!event_loop_ensure_worker(target_loop)) { + /* The timer belongs to the capsule's loop: the message carries that + * loop and the worker dispatches the expiry to it, so a pool loop's + * timers never land in another loop's pending queue. A loop without a + * worker of its own borrows the global shared worker. */ + if (!event_loop_ensure_worker(loop)) { PyErr_SetString(PyExc_RuntimeError, "Event loop has no router or worker"); return NULL; } if (delay_ms < 0) delay_ms = 0; - uint64_t timer_ref_id = atomic_fetch_add(&target_loop->next_callback_id, 1); + uint64_t timer_ref_id = atomic_fetch_add(&g_next_timer_ref, 1); ErlNifEnv *msg_env = enif_alloc_env(); if (msg_env == NULL) { @@ -8072,8 +8141,8 @@ static PyObject *py_schedule_timer_for(PyObject *self, PyObject *args) { return NULL; } - /* Include the target loop resource in message so dispatch goes to correct loop */ - ERL_NIF_TERM loop_term = enif_make_resource(msg_env, target_loop); + /* Include the loop resource in message so dispatch goes to correct loop */ + ERL_NIF_TERM loop_term = enif_make_resource(msg_env, loop); ERL_NIF_TERM msg = enif_make_tuple5( msg_env, @@ -8084,9 +8153,7 @@ static PyObject *py_schedule_timer_for(PyObject *self, PyObject *args) { enif_make_uint64(msg_env, timer_ref_id) ); - /* Use worker_pid when available for scalable I/O */ - ErlNifPid *target_pid = &target_loop->worker_pid; - int send_result = enif_send(NULL, target_pid, msg_env, msg); + int send_result = enif_send(NULL, &loop->worker_pid, msg_env, msg); enif_free_env(msg_env); if (!send_result) { @@ -8535,6 +8602,7 @@ int init_subinterpreter_event_loop(ErlNifEnv *env) { {"set_event_loop_priv_dir", 1, nif_set_event_loop_priv_dir, 0}, \ {"event_loop_new", 0, nif_event_loop_new, 0}, \ {"event_loop_destroy", 1, nif_event_loop_destroy, 0}, \ + {"event_loop_release_python_loop", 1, nif_event_loop_release_python_loop, ERL_NIF_DIRTY_JOB_CPU_BOUND}, \ {"event_loop_set_router", 2, nif_event_loop_set_router, 0}, \ {"event_loop_set_worker", 2, nif_event_loop_set_worker, 0}, \ {"event_loop_set_id", 2, nif_event_loop_set_id, 0}, \ diff --git a/c_src/py_event_loop.h b/c_src/py_event_loop.h index d8adb75..52d890d 100644 --- a/c_src/py_event_loop.h +++ b/c_src/py_event_loop.h @@ -560,6 +560,8 @@ ERL_NIF_TERM nif_event_loop_new(ErlNifEnv *env, int argc, * * NIF: event_loop_destroy(LoopRef) -> ok | {error, Reason} */ +ERL_NIF_TERM nif_event_loop_release_python_loop(ErlNifEnv *env, int argc, + const ERL_NIF_TERM argv[]); ERL_NIF_TERM nif_event_loop_destroy(ErlNifEnv *env, int argc, const ERL_NIF_TERM argv[]); diff --git a/c_src/py_nif.c b/c_src/py_nif.c index 69fc110..54d77d6 100644 --- a/c_src/py_nif.c +++ b/c_src/py_nif.c @@ -260,6 +260,10 @@ static int is_inline_schedule_marker(PyObject *obj); #include "py_exec.c" #include "py_logging.c" #include "py_shared_dict.c" +/* Inline callback path for call requests on worker contexts (defined with + * the context thread below, used by erlang_call_impl). */ +static PyObject *ctx_call_erlang_inline(py_context_t *ctx, const char *func_name, + size_t func_name_len, PyObject *call_args); #include "py_callback.c" #include "py_thread_worker.c" #include "py_event_loop.c" @@ -483,7 +487,11 @@ static void suspended_context_state_destructor(ErlNifEnv *env, void *obj) { enif_release_binary(&state->orig_code); } - /* Release the context resource (was kept in create_suspended_context_state_*) */ + /* Release the env and context resources (kept in create_suspended_context_state_*) */ + if (state->penv != NULL) { + enif_release_resource(state->penv); + state->penv = NULL; + } if (state->ctx != NULL) { enif_release_resource(state->ctx); state->ctx = NULL; @@ -1357,35 +1365,106 @@ static void ctx_queue_cancel_all(py_context_t *ctx) { * ============================================================================ */ /** - * @brief Execute a call request in the OWN_GIL thread + * @brief Set an {error, Reason} response */ -static void ctx_execute_call(py_context_t *ctx) { - /* Decode request from shared_env */ - ERL_NIF_TERM module_term, func_term, args_term, kwargs_term; +static void ctx_set_error(py_context_t *ctx, const char *reason) { + ctx->response_term = enif_make_tuple2(ctx->shared_env, + enif_make_atom(ctx->shared_env, "error"), + enif_make_atom(ctx->shared_env, reason)); + ctx->response_ok = false; +} + +/** + * @brief Resolve Module.Func for a call request or its replay + * + * `__main__` is the caller's namespace: the process-local env when the + * request has one, then the context globals. Other modules go through the + * context module cache. + * + * @return New reference, or NULL with the Python error set + */ +static PyObject *ctx_resolve_call_target(py_context_t *ctx, py_env_resource_t *penv, + const char *module_name, const char *func_name) { + if (strcmp(module_name, "__main__") == 0) { + PyObject *func = NULL; + if (penv != NULL && penv->globals != NULL) { + func = PyDict_GetItemString(penv->globals, func_name); /* Borrowed ref */ + } + if (func == NULL && ctx->globals != NULL) { + func = PyDict_GetItemString(ctx->globals, func_name); /* Borrowed ref */ + } + if (func != NULL) { + Py_INCREF(func); + return func; + } + } + PyObject *module = context_get_module(ctx, module_name); /* Borrowed ref (cached) */ + if (module == NULL) { + return NULL; + } + return PyObject_GetAttrString(module, func_name); +} + +/** + * @brief Build the positional argument tuple of a call request + * + * @return New tuple, or NULL with *reason set to the error atom name + */ +static PyObject *ctx_build_call_args(ErlNifEnv *env, ERL_NIF_TERM args_term, + const char **reason) { + unsigned int args_len; + if (!enif_get_list_length(env, args_term, &args_len)) { + *reason = "invalid_args"; + return NULL; + } + PyObject *args = PyTuple_New(args_len); + if (args == NULL) { + PyErr_Clear(); + *reason = "arg_conversion_failed"; + return NULL; + } + ERL_NIF_TERM head, tail = args_term; + for (unsigned int i = 0; i < args_len; i++) { + enif_get_list_cell(env, tail, &head, &tail); + PyObject *arg = term_to_py(env, head); + if (arg == NULL) { + PyErr_Clear(); + Py_DECREF(args); + *reason = "arg_conversion_failed"; + return NULL; + } + PyTuple_SET_ITEM(args, i, arg); + } + return args; +} + +/** + * @brief Execute a call request, with or without a process-local env + * + * Shared body of ctx_execute_call and ctx_execute_call_with_env. An + * erlang.call in the function waits inline (ctx_call_erlang_inline): the + * thread serves this context's queue meanwhile and nothing is replayed. + * + * @param penv Process-local env for `__main__` lookups, or NULL + */ +static void ctx_execute_call_in(py_context_t *ctx, py_env_resource_t *penv) { + /* Decode request from shared_env: {Module, Func, Args, Kwargs} */ const ERL_NIF_TERM *tuple_terms; int tuple_arity; if (!enif_get_tuple(ctx->shared_env, ctx->request_term, &tuple_arity, &tuple_terms) || tuple_arity < 4) { - ctx->response_term = enif_make_tuple2(ctx->shared_env, - enif_make_atom(ctx->shared_env, "error"), - enif_make_atom(ctx->shared_env, "invalid_request")); - ctx->response_ok = false; + ctx_set_error(ctx, "invalid_request"); return; } - module_term = tuple_terms[0]; - func_term = tuple_terms[1]; - args_term = tuple_terms[2]; - kwargs_term = tuple_terms[3]; + ERL_NIF_TERM args_term = tuple_terms[2]; + ERL_NIF_TERM kwargs_term = tuple_terms[3]; ErlNifBinary module_bin, func_bin; - if (!enif_inspect_binary(ctx->shared_env, module_term, &module_bin) || - !enif_inspect_binary(ctx->shared_env, func_term, &func_bin)) { - ctx->response_term = enif_make_tuple2(ctx->shared_env, - enif_make_atom(ctx->shared_env, "error"), - enif_make_atom(ctx->shared_env, "invalid_module_or_func")); - ctx->response_ok = false; + if (!enif_inspect_binary(ctx->shared_env, tuple_terms[0], &module_bin) || + !enif_inspect_binary(ctx->shared_env, tuple_terms[1], &func_bin)) { + ctx_set_error(ctx, "invalid_module_or_func"); return; } @@ -1395,83 +1474,33 @@ static void ctx_execute_call(py_context_t *ctx) { if (module_name == NULL || func_name_str == NULL) { enif_free(module_name); enif_free(func_name_str); - ctx->response_term = enif_make_tuple2(ctx->shared_env, - enif_make_atom(ctx->shared_env, "error"), - enif_make_atom(ctx->shared_env, "alloc_failed")); - ctx->response_ok = false; + ctx_set_error(ctx, "alloc_failed"); return; } - PyObject *module = NULL; - PyObject *func = NULL; - - /* Special handling for __main__ module - check ctx->globals first */ - if (strcmp(module_name, "__main__") == 0) { - func = PyDict_GetItemString(ctx->globals, func_name_str); /* Borrowed ref */ - if (func != NULL) { - Py_INCREF(func); - } - } - - if (func == NULL) { - /* Get or import module */ - module = context_get_module(ctx, module_name); - if (module == NULL) { - ctx->response_term = make_py_error(ctx->shared_env); - ctx->response_ok = false; - enif_free(module_name); - enif_free(func_name_str); - return; - } - - /* Get function */ - func = PyObject_GetAttrString(module, func_name_str); - if (func == NULL) { - ctx->response_term = make_py_error(ctx->shared_env); - ctx->response_ok = false; - enif_free(module_name); - enif_free(func_name_str); - return; - } - } + /* Thread-local state for callbacks: with tl_current_context set and + * suspension left off, an erlang.call in the function takes the inline + * path (ctx_call_erlang_inline) and the function runs exactly once. */ + py_context_t *prev_context = tl_current_context; + tl_current_context = ctx; + py_env_resource_t *prev_local_env = tl_current_local_env; + tl_current_local_env = penv; + PyObject *func = ctx_resolve_call_target(ctx, penv, module_name, func_name_str); enif_free(module_name); enif_free(func_name_str); - - /* Convert args */ - unsigned int args_len; - if (!enif_get_list_length(ctx->shared_env, args_term, &args_len)) { - Py_DECREF(func); - ctx->response_term = enif_make_tuple2(ctx->shared_env, - enif_make_atom(ctx->shared_env, "error"), - enif_make_atom(ctx->shared_env, "invalid_args")); + if (func == NULL) { + ctx->response_term = make_py_error(ctx->shared_env); ctx->response_ok = false; - return; + goto cleanup; } - PyObject *args = PyTuple_New(args_len); + const char *reason = NULL; + PyObject *args = ctx_build_call_args(ctx->shared_env, args_term, &reason); if (args == NULL) { Py_DECREF(func); - ctx->response_term = enif_make_tuple2(ctx->shared_env, - enif_make_atom(ctx->shared_env, "error"), - enif_make_atom(ctx->shared_env, "arg_conversion_failed")); - ctx->response_ok = false; - return; - } - ERL_NIF_TERM head, tail = args_term; - for (unsigned int i = 0; i < args_len; i++) { - enif_get_list_cell(ctx->shared_env, tail, &head, &tail); - PyObject *arg = term_to_py(ctx->shared_env, head); - if (arg == NULL) { - Py_DECREF(args); - Py_DECREF(func); - ctx->response_term = enif_make_tuple2(ctx->shared_env, - enif_make_atom(ctx->shared_env, "error"), - enif_make_atom(ctx->shared_env, "arg_conversion_failed")); - ctx->response_ok = false; - return; - } - PyTuple_SET_ITEM(args, i, arg); + ctx_set_error(ctx, reason); + goto cleanup; } /* Convert kwargs */ @@ -1486,16 +1515,30 @@ static void ctx_execute_call(py_context_t *ctx) { Py_DECREF(args); Py_XDECREF(kwargs); - if (py_result == NULL) { - ctx->response_term = make_py_error(ctx->shared_env); - ctx->response_ok = false; - } else { + if (py_result != NULL) { ERL_NIF_TERM term_result = py_to_term(ctx->shared_env, py_result); Py_DECREF(py_result); ctx->response_term = enif_make_tuple2(ctx->shared_env, enif_make_atom(ctx->shared_env, "ok"), term_result); ctx->response_ok = true; + } else { + ctx->response_term = make_py_error(ctx->shared_env); + ctx->response_ok = false; } + +cleanup: + /* Restore thread-local state */ + tl_current_local_env = prev_local_env; + tl_current_context = prev_context; +} + +/** + * @brief Execute a call request in the OWN_GIL thread + * + * `__main__` functions are looked up in ctx->globals. + */ +static void ctx_execute_call(py_context_t *ctx) { + ctx_execute_call_in(ctx, NULL); } /** @@ -1871,7 +1914,7 @@ static void ctx_execute_eval_with_env(py_context_t *ctx) { PyErr_Clear(); /* Create suspended state for callback handling */ suspended_context_state_t *suspended = create_suspended_context_state_for_eval( - ctx->shared_env, ctx, &code_bin, tuple_terms[1]); + ctx->shared_env, ctx, &code_bin, tuple_terms[1], penv); if (suspended == NULL) { tl_pending_callback = false; Py_CLEAR(tl_pending_args); @@ -1974,7 +2017,7 @@ static void ctx_execute_eval_with_env(py_context_t *ctx) { if (tl_pending_callback) { PyErr_Clear(); suspended_context_state_t *suspended = create_suspended_context_state_for_eval( - ctx->shared_env, ctx, &code_bin, tuple_terms[1]); + ctx->shared_env, ctx, &code_bin, tuple_terms[1], penv); if (suspended == NULL) { tl_pending_callback = false; Py_CLEAR(tl_pending_args); @@ -2043,10 +2086,7 @@ static void ctx_execute_call_with_env(py_context_t *ctx) { ctx->local_env_ptr = NULL; /* Clear after use */ if (penv == NULL || penv->globals == NULL) { - ctx->response_term = enif_make_tuple2(ctx->shared_env, - enif_make_atom(ctx->shared_env, "error"), - enif_make_atom(ctx->shared_env, "invalid_env")); - ctx->response_ok = false; + ctx_set_error(ctx, "invalid_env"); return; } @@ -2054,159 +2094,11 @@ static void ctx_execute_call_with_env(py_context_t *ctx) { * Compare env's interp_id with the current Python interpreter's ID. */ PyInterpreterState *current_interp = PyInterpreterState_Get(); if (current_interp != NULL && penv->interp_id != PyInterpreterState_GetID(current_interp)) { - ctx->response_term = enif_make_tuple2(ctx->shared_env, - enif_make_atom(ctx->shared_env, "error"), - enif_make_atom(ctx->shared_env, "env_wrong_interpreter")); - ctx->response_ok = false; - return; - } - - /* Decode request from shared_env: {Module, Func, Args, Kwargs} */ - ERL_NIF_TERM module_term, func_term, args_term, kwargs_term; - const ERL_NIF_TERM *tuple_terms; - int tuple_arity; - - if (!enif_get_tuple(ctx->shared_env, ctx->request_term, &tuple_arity, &tuple_terms) || - tuple_arity < 4) { - ctx->response_term = enif_make_tuple2(ctx->shared_env, - enif_make_atom(ctx->shared_env, "error"), - enif_make_atom(ctx->shared_env, "invalid_request")); - ctx->response_ok = false; - return; - } - - module_term = tuple_terms[0]; - func_term = tuple_terms[1]; - args_term = tuple_terms[2]; - kwargs_term = tuple_terms[3]; - - ErlNifBinary module_bin, func_bin; - if (!enif_inspect_binary(ctx->shared_env, module_term, &module_bin) || - !enif_inspect_binary(ctx->shared_env, func_term, &func_bin)) { - ctx->response_term = enif_make_tuple2(ctx->shared_env, - enif_make_atom(ctx->shared_env, "error"), - enif_make_atom(ctx->shared_env, "invalid_module_or_func")); - ctx->response_ok = false; - return; - } - - char *module_name = binary_to_string(&module_bin); - char *func_name_str = binary_to_string(&func_bin); - - if (module_name == NULL || func_name_str == NULL) { - enif_free(module_name); - enif_free(func_name_str); - ctx->response_term = enif_make_tuple2(ctx->shared_env, - enif_make_atom(ctx->shared_env, "error"), - enif_make_atom(ctx->shared_env, "alloc_failed")); - ctx->response_ok = false; + ctx_set_error(ctx, "env_wrong_interpreter"); return; } - /* Set thread-local env for callback support */ - py_env_resource_t *prev_local_env = tl_current_local_env; - tl_current_local_env = penv; - - PyObject *func = NULL; - - /* Special handling for __main__ module - look up in process-local globals */ - if (strcmp(module_name, "__main__") == 0) { - func = PyDict_GetItemString(penv->globals, func_name_str); /* Borrowed ref */ - if (func != NULL) { - Py_INCREF(func); - } - } - - if (func == NULL) { - /* Get or import module from context cache */ - PyObject *module = context_get_module(ctx, module_name); - if (module == NULL) { - enif_free(module_name); - enif_free(func_name_str); - tl_current_local_env = prev_local_env; - ctx->response_term = make_py_error(ctx->shared_env); - ctx->response_ok = false; - return; - } - - /* Get function */ - func = PyObject_GetAttrString(module, func_name_str); - if (func == NULL) { - enif_free(module_name); - enif_free(func_name_str); - tl_current_local_env = prev_local_env; - ctx->response_term = make_py_error(ctx->shared_env); - ctx->response_ok = false; - return; - } - } - - enif_free(module_name); - enif_free(func_name_str); - - /* Convert args */ - unsigned int args_len; - if (!enif_get_list_length(ctx->shared_env, args_term, &args_len)) { - Py_DECREF(func); - tl_current_local_env = prev_local_env; - ctx->response_term = enif_make_tuple2(ctx->shared_env, - enif_make_atom(ctx->shared_env, "error"), - enif_make_atom(ctx->shared_env, "invalid_args")); - ctx->response_ok = false; - return; - } - - PyObject *args = PyTuple_New(args_len); - if (args == NULL) { - Py_DECREF(func); - tl_current_local_env = prev_local_env; - ctx->response_term = enif_make_tuple2(ctx->shared_env, - enif_make_atom(ctx->shared_env, "error"), - enif_make_atom(ctx->shared_env, "arg_conversion_failed")); - ctx->response_ok = false; - return; - } - ERL_NIF_TERM head, tail = args_term; - for (unsigned int i = 0; i < args_len; i++) { - enif_get_list_cell(ctx->shared_env, tail, &head, &tail); - PyObject *arg = term_to_py(ctx->shared_env, head); - if (arg == NULL) { - Py_DECREF(args); - Py_DECREF(func); - tl_current_local_env = prev_local_env; - ctx->response_term = enif_make_tuple2(ctx->shared_env, - enif_make_atom(ctx->shared_env, "error"), - enif_make_atom(ctx->shared_env, "arg_conversion_failed")); - ctx->response_ok = false; - return; - } - PyTuple_SET_ITEM(args, i, arg); - } - - /* Convert kwargs */ - PyObject *kwargs = NULL; - if (enif_is_map(ctx->shared_env, kwargs_term)) { - kwargs = term_to_py(ctx->shared_env, kwargs_term); - } - - /* Call the function */ - PyObject *py_result = PyObject_Call(func, args, kwargs); - Py_DECREF(func); - Py_DECREF(args); - Py_XDECREF(kwargs); - - tl_current_local_env = prev_local_env; - - if (py_result == NULL) { - ctx->response_term = make_py_error(ctx->shared_env); - ctx->response_ok = false; - } else { - ERL_NIF_TERM term_result = py_to_term(ctx->shared_env, py_result); - Py_DECREF(py_result); - ctx->response_term = enif_make_tuple2(ctx->shared_env, - enif_make_atom(ctx->shared_env, "ok"), term_result); - ctx->response_ok = true; - } + ctx_execute_call_in(ctx, penv); } /** @@ -2459,6 +2351,235 @@ static void ctx_execute_request(py_context_t *ctx) { * requiring subinterpreter support. * ============================================================================ */ +/** + * @brief Serve one dequeued request on the context thread + * + * Executes the request through the context's mirror fields and delivers + * the result (message in async mode, condvar otherwise). Called by the + * request loop and, with @p nested set, by ctx_call_erlang_inline while an + * erlang.call waits for its callback: the GIL is then already held and the + * caller saves and restores the mirror fields of the outer request. + */ +static void ctx_serve_request(py_context_t *ctx, ctx_request_t *req, bool nested) { + /* Check if request was cancelled while queued */ + if (atomic_load(&req->cancelled)) { + /* Request cancelled - deliver error without processing */ + if (req->async_mode) { + /* Async mode: send cancellation message */ + enif_clear_env(ctx->msg_env); + ERL_NIF_TERM cancel_msg = enif_make_tuple3(ctx->msg_env, + enif_make_atom(ctx->msg_env, "py_result"), + enif_make_copy(ctx->msg_env, req->request_id), + enif_make_tuple2(ctx->msg_env, + enif_make_atom(ctx->msg_env, "error"), + enif_make_atom(ctx->msg_env, "cancelled"))); + enif_send(NULL, &req->caller_pid, ctx->msg_env, cancel_msg); + } else { + /* Blocking mode: signal condvar */ + req->result_env = enif_alloc_env(); + if (req->result_env) { + req->result = enif_make_tuple2(req->result_env, + enif_make_atom(req->result_env, "error"), + enif_make_atom(req->result_env, "cancelled")); + } + req->success = false; + + pthread_mutex_lock(&req->mutex); + atomic_store(&req->completed, true); + pthread_cond_signal(&req->cond); + pthread_mutex_unlock(&req->mutex); + } + + ctx_request_release(req); + return; + } + + /* Populate legacy compatibility fields from request */ + ctx->shared_env = req->request_env; + ctx->request_type = req->type; + ctx->request_term = req->request_data; + ctx->reactor_buffer_ptr = req->reactor_buffer_ptr; + ctx->local_env_ptr = req->local_env_ptr; + ctx->has_current_caller = req->async_mode; + if (req->async_mode) { + ctx->current_caller = req->caller_pid; + } + ctx->response_ok = false; + ctx->response_term = 0; + + if (nested) { + /* Inside an outer request: GIL held, exec marker already set */ + ctx_execute_request(ctx); + } else { + /* Acquire GIL and process the request. + * exec_enter before / exec_leave after the GIL (see the locking + * invariant on py_context::interrupt_mutex). */ + py_context_exec_enter(ctx); + PyGILState_STATE gstate = PyGILState_Ensure(); + ctx_execute_request(ctx); /* Reuse execute functions */ + PyGILState_Release(gstate); + py_context_exec_leave(ctx); + } + + /* Copy response to request struct */ + req->result_env = enif_alloc_env(); + if (req->result_env && ctx->response_term != 0) { + req->result = enif_make_copy(req->result_env, ctx->response_term); + } else if (req->result_env) { + req->result = enif_make_tuple2(req->result_env, + enif_make_atom(req->result_env, "error"), + enif_make_atom(req->result_env, "no_response")); + } + req->success = ctx->response_ok; + + /* Clear legacy fields */ + ctx->shared_env = NULL; + ctx->request_type = CTX_REQ_NONE; + ctx->request_term = 0; + ctx->reactor_buffer_ptr = NULL; + ctx->local_env_ptr = NULL; + ctx->has_current_caller = false; + + /* Deliver result - async or blocking */ + if (req->async_mode) { + /* Async mode: send result message to caller */ + enif_clear_env(ctx->msg_env); + ERL_NIF_TERM result_msg = enif_make_tuple3(ctx->msg_env, + enif_make_atom(ctx->msg_env, "py_result"), + enif_make_copy(ctx->msg_env, req->request_id), + req->result_env ? enif_make_copy(ctx->msg_env, req->result) + : enif_make_tuple2(ctx->msg_env, + enif_make_atom(ctx->msg_env, "error"), + enif_make_atom(ctx->msg_env, "no_result"))); + enif_send(NULL, &req->caller_pid, ctx->msg_env, result_msg); + } else { + /* Blocking mode: signal condvar */ + pthread_mutex_lock(&req->mutex); + atomic_store(&req->completed, true); + pthread_cond_signal(&req->cond); + pthread_mutex_unlock(&req->mutex); + } + + /* Release queue's reference to request */ + ctx_request_release(req); +} + +/** + * @brief erlang.call from a call request: wait inline, serving the queue + * + * Sends {py_callback, CallbackId, Name, Args} to the py_context process + * that issued the current request, then waits with the GIL released for + * the reply delivered by context_callback_reply. While waiting, every + * request queued for this context is served on this thread, so the + * callback's nested py:call (bound to this context by py_context) runs + * here, at any depth, and the outer Python frame is never replayed. + * + * @return New reference to the callback result, or NULL with a Python + * error set (callback error, context shutting down) + */ +static PyObject *ctx_call_erlang_inline(py_context_t *ctx, const char *func_name, + size_t func_name_len, PyObject *call_args) { + uint64_t callback_id = atomic_fetch_add(&g_callback_id_counter, 1); + + ErlNifEnv *msg_env = enif_alloc_env(); + if (msg_env == NULL) { + PyErr_SetString(PyExc_MemoryError, "Failed to allocate callback env"); + return NULL; + } + ERL_NIF_TERM name_term; + unsigned char *name_buf = enif_make_new_binary(msg_env, func_name_len, &name_term); + memcpy(name_buf, func_name, func_name_len); + ERL_NIF_TERM args_term = py_to_term(msg_env, call_args); + ERL_NIF_TERM msg = enif_make_tuple4(msg_env, + enif_make_atom(msg_env, "py_callback"), + enif_make_uint64(msg_env, callback_id), + name_term, args_term); + int sent = enif_send(NULL, &ctx->current_caller, msg_env, msg); + enif_free_env(msg_env); + if (!sent) { + PyErr_SetString(PyExc_RuntimeError, "erlang.call: context process is gone"); + return NULL; + } + + /* Save the outer request's mirror fields: nested requests reuse them */ + ErlNifEnv *saved_shared_env = ctx->shared_env; + int saved_request_type = ctx->request_type; + ERL_NIF_TERM saved_request_term = ctx->request_term; + ERL_NIF_TERM saved_response_term = ctx->response_term; + bool saved_response_ok = ctx->response_ok; + void *saved_reactor_buffer_ptr = ctx->reactor_buffer_ptr; + void *saved_local_env_ptr = ctx->local_env_ptr; + ErlNifPid saved_caller = ctx->current_caller; + bool saved_has_caller = ctx->has_current_caller; + py_env_resource_t *saved_local_env = tl_current_local_env; + + unsigned char *reply = NULL; + size_t reply_len = 0; + bool shutdown = false; + + PyThreadState *tstate = PyEval_SaveThread(); + for (;;) { + pthread_mutex_lock(&ctx->queue_mutex); + while (!(ctx->cb_reply_ready && ctx->cb_reply_id == callback_id) && + ctx->queue_head == NULL && + !atomic_load(&ctx->shutdown_requested)) { + pthread_cond_wait(&ctx->queue_not_empty, &ctx->queue_mutex); + } + if (ctx->cb_reply_ready && ctx->cb_reply_id == callback_id) { + reply = ctx->cb_reply_data; + reply_len = ctx->cb_reply_len; + ctx->cb_reply_data = NULL; + ctx->cb_reply_len = 0; + ctx->cb_reply_ready = false; + pthread_mutex_unlock(&ctx->queue_mutex); + break; + } + if (atomic_load(&ctx->shutdown_requested)) { + pthread_mutex_unlock(&ctx->queue_mutex); + shutdown = true; + break; + } + ctx_request_t *req = ctx->queue_head; + if (req->type == CTX_REQ_SHUTDOWN) { + /* Leave the sentinel for the request loop; unwind this request */ + pthread_mutex_unlock(&ctx->queue_mutex); + shutdown = true; + break; + } + ctx->queue_head = req->next; + if (ctx->queue_head == NULL) { + ctx->queue_tail = NULL; + } + req->next = NULL; + pthread_mutex_unlock(&ctx->queue_mutex); + + PyEval_RestoreThread(tstate); + ctx_serve_request(ctx, req, true); + tstate = PyEval_SaveThread(); + } + PyEval_RestoreThread(tstate); + + /* Restore the outer request */ + ctx->shared_env = saved_shared_env; + ctx->request_type = saved_request_type; + ctx->request_term = saved_request_term; + ctx->response_term = saved_response_term; + ctx->response_ok = saved_response_ok; + ctx->reactor_buffer_ptr = saved_reactor_buffer_ptr; + ctx->local_env_ptr = saved_local_env_ptr; + ctx->current_caller = saved_caller; + ctx->has_current_caller = saved_has_caller; + tl_current_local_env = saved_local_env; + + if (shutdown) { + PyErr_SetString(PyExc_RuntimeError, "erlang.call: context shut down while waiting"); + return NULL; + } + PyObject *result = parse_callback_response(reply, reply_len); + enif_free(reply); + return result; +} + /** * @brief Main loop for worker context thread (main interpreter mode) * @@ -2527,97 +2648,7 @@ static void *ctx_thread_main_worker(void *arg) { break; } - /* Check if request was cancelled while queued */ - if (atomic_load(&req->cancelled)) { - /* Request cancelled - deliver error without processing */ - if (req->async_mode) { - /* Async mode: send cancellation message */ - enif_clear_env(ctx->msg_env); - ERL_NIF_TERM cancel_msg = enif_make_tuple3(ctx->msg_env, - enif_make_atom(ctx->msg_env, "py_result"), - enif_make_copy(ctx->msg_env, req->request_id), - enif_make_tuple2(ctx->msg_env, - enif_make_atom(ctx->msg_env, "error"), - enif_make_atom(ctx->msg_env, "cancelled"))); - enif_send(NULL, &req->caller_pid, ctx->msg_env, cancel_msg); - } else { - /* Blocking mode: signal condvar */ - req->result_env = enif_alloc_env(); - if (req->result_env) { - req->result = enif_make_tuple2(req->result_env, - enif_make_atom(req->result_env, "error"), - enif_make_atom(req->result_env, "cancelled")); - } - req->success = false; - - pthread_mutex_lock(&req->mutex); - atomic_store(&req->completed, true); - pthread_cond_signal(&req->cond); - pthread_mutex_unlock(&req->mutex); - } - - ctx_request_release(req); - continue; - } - - /* Populate legacy compatibility fields from request */ - ctx->shared_env = req->request_env; - ctx->request_type = req->type; - ctx->request_term = req->request_data; - ctx->reactor_buffer_ptr = req->reactor_buffer_ptr; - ctx->local_env_ptr = req->local_env_ptr; - ctx->response_ok = false; - ctx->response_term = 0; - - /* Acquire GIL and process the request. - * exec_enter before / exec_leave after the GIL (see the locking - * invariant on py_context::interrupt_mutex). */ - py_context_exec_enter(ctx); - gstate = PyGILState_Ensure(); - ctx_execute_request(ctx); /* Reuse execute functions */ - PyGILState_Release(gstate); - py_context_exec_leave(ctx); - - /* Copy response to request struct */ - req->result_env = enif_alloc_env(); - if (req->result_env && ctx->response_term != 0) { - req->result = enif_make_copy(req->result_env, ctx->response_term); - } else if (req->result_env) { - req->result = enif_make_tuple2(req->result_env, - enif_make_atom(req->result_env, "error"), - enif_make_atom(req->result_env, "no_response")); - } - req->success = ctx->response_ok; - - /* Clear legacy fields */ - ctx->shared_env = NULL; - ctx->request_type = CTX_REQ_NONE; - ctx->request_term = 0; - ctx->reactor_buffer_ptr = NULL; - ctx->local_env_ptr = NULL; - - /* Deliver result - async or blocking */ - if (req->async_mode) { - /* Async mode: send result message to caller */ - enif_clear_env(ctx->msg_env); - ERL_NIF_TERM result_msg = enif_make_tuple3(ctx->msg_env, - enif_make_atom(ctx->msg_env, "py_result"), - enif_make_copy(ctx->msg_env, req->request_id), - req->result_env ? enif_make_copy(ctx->msg_env, req->result) - : enif_make_tuple2(ctx->msg_env, - enif_make_atom(ctx->msg_env, "error"), - enif_make_atom(ctx->msg_env, "no_result"))); - enif_send(NULL, &req->caller_pid, ctx->msg_env, result_msg); - } else { - /* Blocking mode: signal condvar */ - pthread_mutex_lock(&req->mutex); - atomic_store(&req->completed, true); - pthread_cond_signal(&req->cond); - pthread_mutex_unlock(&req->mutex); - } - - /* Release queue's reference to request */ - ctx_request_release(req); + ctx_serve_request(ctx, req, false); } /* Cleanup: release namespace dictionaries under GIL */ @@ -2798,6 +2829,11 @@ static void ctx_thread_shutdown_worker(py_context_t *ctx) { enif_free_env(ctx->msg_env); ctx->msg_env = NULL; } + if (ctx->cb_reply_data != NULL) { + enif_free(ctx->cb_reply_data); + ctx->cb_reply_data = NULL; + ctx->cb_reply_ready = false; + } pthread_cond_destroy(&ctx->queue_not_empty); pthread_mutex_destroy(&ctx->queue_mutex); @@ -3433,6 +3469,11 @@ static void ctx_thread_shutdown_owngil(py_context_t *ctx) { enif_free_env(ctx->msg_env); ctx->msg_env = NULL; } + if (ctx->cb_reply_data != NULL) { + enif_free(ctx->cb_reply_data); + ctx->cb_reply_data = NULL; + ctx->cb_reply_ready = false; + } pthread_cond_destroy(&ctx->queue_not_empty); pthread_mutex_destroy(&ctx->queue_mutex); @@ -3504,6 +3545,11 @@ static ERL_NIF_TERM nif_context_create(ErlNifEnv *env, int argc, const ERL_NIF_T ctx->has_callback_handler = false; ctx->callback_pipe[0] = -1; ctx->callback_pipe[1] = -1; + ctx->has_current_caller = false; + ctx->cb_reply_id = 0; + ctx->cb_reply_data = NULL; + ctx->cb_reply_len = 0; + ctx->cb_reply_ready = false; ctx->globals = NULL; ctx->locals = NULL; ctx->module_cache = NULL; @@ -4689,6 +4735,49 @@ static ERL_NIF_TERM nif_context_write_callback_response(ErlNifEnv *env, int argc * A proper implementation would add PY_CMD_RESUME and dispatch to the * dedicated thread. */ +/** + * nif_context_callback_reply(ContextRef, CallbackId, ResultBinary) -> ok | {error, Reason} + * + * Delivers the result of an inline callback (see ctx_call_erlang_inline) + * to the context thread waiting for it. ResultBinary is the frame + * parse_callback_response reads: a status byte then the payload. + */ +static ERL_NIF_TERM nif_context_callback_reply(ErlNifEnv *env, int argc, const ERL_NIF_TERM argv[]) { + (void)argc; + py_context_t *ctx; + ErlNifUInt64 callback_id; + ErlNifBinary result_bin; + + if (!enif_get_resource(env, argv[0], PY_CONTEXT_RESOURCE_TYPE, (void **)&ctx)) { + return make_error(env, "invalid_context"); + } + if (!enif_get_uint64(env, argv[1], &callback_id)) { + return make_error(env, "invalid_callback_id"); + } + if (!enif_inspect_binary(env, argv[2], &result_bin)) { + return make_error(env, "invalid_result"); + } + + unsigned char *data = enif_alloc(result_bin.size > 0 ? result_bin.size : 1); + if (data == NULL) { + return make_error(env, "alloc_failed"); + } + memcpy(data, result_bin.data, result_bin.size); + + pthread_mutex_lock(&ctx->queue_mutex); + if (ctx->cb_reply_data != NULL) { + enif_free(ctx->cb_reply_data); /* A reply nobody waited for */ + } + ctx->cb_reply_id = callback_id; + ctx->cb_reply_data = data; + ctx->cb_reply_len = result_bin.size; + ctx->cb_reply_ready = true; + pthread_cond_broadcast(&ctx->queue_not_empty); + pthread_mutex_unlock(&ctx->queue_mutex); + + return ATOM_OK; +} + static ERL_NIF_TERM nif_context_resume(ErlNifEnv *env, int argc, const ERL_NIF_TERM argv[]) { (void)argc; py_context_t *ctx; @@ -4755,6 +4844,10 @@ static ERL_NIF_TERM nif_context_resume(ErlNifEnv *env, int argc, const ERL_NIF_T suspended_context_state_t *prev_suspended = tl_current_context_suspended; tl_current_context_suspended = state; + /* Replay in the namespace of the original request */ + py_env_resource_t *prev_local_env = tl_current_local_env; + tl_current_local_env = state->penv; + /* Reset callback result index for this replay */ state->callback_result_index = 0; @@ -4775,17 +4868,8 @@ static ERL_NIF_TERM nif_context_resume(ErlNifEnv *env, int argc, const ERL_NIF_T memcpy(func_name, state->orig_func.data, state->orig_func.size); func_name[state->orig_func.size] = '\0'; - /* Get the function */ - PyObject *func = NULL; - PyObject *module = context_get_module(ctx, module_name); - if (module == NULL) { - enif_free(module_name); - enif_free(func_name); - result = make_py_error(env); - goto cleanup; - } - - func = PyObject_GetAttrString(module, func_name); + /* Get the function, in the namespace the original request used */ + PyObject *func = ctx_resolve_call_target(ctx, state->penv, module_name, func_name); if (func == NULL) { enif_free(module_name); enif_free(func_name); @@ -4848,7 +4932,7 @@ static ERL_NIF_TERM nif_context_resume(ErlNifEnv *env, int argc, const ERL_NIF_T /* Create new suspended context state for nested callback */ suspended_context_state_t *nested = create_suspended_context_state_for_call( env, ctx, &state->orig_module, &state->orig_func, - state->orig_args, state->orig_kwargs); + state->orig_args, state->orig_kwargs, state->penv); if (nested == NULL) { tl_pending_callback = false; @@ -4884,17 +4968,37 @@ static ERL_NIF_TERM nif_context_resume(ErlNifEnv *env, int argc, const ERL_NIF_T memcpy(code, state->orig_code.data, state->orig_code.size); code[state->orig_code.size] = '\0'; - /* Update locals if provided */ - if (enif_is_map(state->orig_env, state->orig_locals)) { - PyObject *new_locals = term_to_py(state->orig_env, state->orig_locals); - if (new_locals != NULL && PyDict_Check(new_locals)) { - PyDict_Update(ctx->locals, new_locals); - Py_DECREF(new_locals); + PyObject *py_result; + if (state->penv != NULL && state->penv->globals != NULL) { + /* Process-local env: same namespaces as ctx_execute_eval_with_env */ + PyObject *eval_locals = PyDict_Copy(state->penv->globals); + if (eval_locals == NULL) { + enif_free(code); + result = make_py_error(env); + goto cleanup; + } + if (enif_is_map(state->orig_env, state->orig_locals)) { + PyObject *locals_map = term_to_py(state->orig_env, state->orig_locals); + if (locals_map != NULL && PyDict_Check(locals_map)) { + PyDict_Merge(eval_locals, locals_map, 1); + } + Py_XDECREF(locals_map); + } + py_result = PyRun_String(code, Py_eval_input, state->penv->globals, eval_locals); + Py_DECREF(eval_locals); + } else { + /* Update locals if provided */ + if (enif_is_map(state->orig_env, state->orig_locals)) { + PyObject *new_locals = term_to_py(state->orig_env, state->orig_locals); + if (new_locals != NULL && PyDict_Check(new_locals)) { + PyDict_Update(ctx->locals, new_locals); + Py_DECREF(new_locals); + } } - } - /* Compile and evaluate (replay with cached result) */ - PyObject *py_result = PyRun_String(code, Py_eval_input, ctx->globals, ctx->locals); + /* Compile and evaluate (replay with cached result) */ + py_result = PyRun_String(code, Py_eval_input, ctx->globals, ctx->locals); + } enif_free(code); if (py_result == NULL) { @@ -4904,7 +5008,7 @@ static ERL_NIF_TERM nif_context_resume(ErlNifEnv *env, int argc, const ERL_NIF_T /* Create new suspended context state for nested callback */ suspended_context_state_t *nested = create_suspended_context_state_for_eval( - env, ctx, &state->orig_code, state->orig_locals); + env, ctx, &state->orig_code, state->orig_locals, state->penv); if (nested == NULL) { tl_pending_callback = false; @@ -4936,6 +5040,7 @@ static ERL_NIF_TERM nif_context_resume(ErlNifEnv *env, int argc, const ERL_NIF_T cleanup: /* Restore thread-local state */ + tl_current_local_env = prev_local_env; tl_current_context_suspended = prev_suspended; tl_allow_suspension = prev_allow_suspension; tl_current_context = prev_context; @@ -5939,6 +6044,7 @@ static ErlNifFunc nif_funcs[] = { {"context_get_callback_pipe", 1, nif_context_get_callback_pipe, 0}, {"context_write_callback_response", 2, nif_context_write_callback_response, ERL_NIF_DIRTY_JOB_IO_BOUND}, {"context_resume", 3, nif_context_resume, ERL_NIF_DIRTY_JOB_CPU_BOUND}, + {"context_callback_reply", 3, nif_context_callback_reply, 0}, {"context_cancel_resume", 2, nif_context_cancel_resume, 0}, {"ref_wrap", 2, nif_ref_wrap, 0}, {"is_ref", 1, nif_is_ref, 0}, diff --git a/c_src/py_nif.h b/c_src/py_nif.h index 5dab26c..be393c0 100644 --- a/c_src/py_nif.h +++ b/c_src/py_nif.h @@ -790,6 +790,32 @@ struct py_context { /** @brief Process-local env pointer for current request */ void *local_env_ptr; + /** @brief py_context process that issued the current request (async + * mode); it runs the callbacks of an inline erlang.call */ + ErlNifPid current_caller; + + /** @brief True when current_caller is valid for the request in flight */ + bool has_current_caller; + + /* ========== Inline callback reply (queue_mutex) ========== */ + /* An erlang.call from a call request blocks the context thread, which + * serves its own queue meanwhile; py_context delivers the callback's + * result here through context_callback_reply and signals + * queue_not_empty. One slot: callbacks nest strictly, so the innermost + * waiter is the only one waiting. */ + + /** @brief Callback id the pending reply belongs to */ + uint64_t cb_reply_id; + + /** @brief Reply frame (status byte + payload), owned here until consumed */ + unsigned char *cb_reply_data; + + /** @brief Length of cb_reply_data */ + size_t cb_reply_len; + + /** @brief True while a reply waits in the slot */ + bool cb_reply_ready; + #ifdef HAVE_SUBINTERPRETERS /* ========== OWN_GIL specific fields ========== */ @@ -1075,6 +1101,10 @@ typedef struct { /** @brief Context for replay */ py_context_t *ctx; + /** @brief Process-local env of the original request, or NULL. + * Kept alive so the replay runs in the caller's namespace. */ + struct py_env_resource *penv; + /** @brief Unique identifier for this callback */ uint64_t callback_id; @@ -1254,7 +1284,7 @@ extern ErlNifResourceType *INLINE_CONTINUATION_RESOURCE_TYPE; * Erlang GC drops the reference, triggering the destructor which frees * the Python dicts. */ -typedef struct { +typedef struct py_env_resource { /** @brief Global namespace dictionary */ PyObject *globals; /** @brief Local namespace dictionary (same as globals for module-level execution) */ diff --git a/docs/architecture.md b/docs/architecture.md index 3ae029f..f95a020 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -107,17 +107,24 @@ docstring of `priv/_erlang_impl/_isolated.py`. `erlang_call_impl` in `c_src/py_callback.c` chooses one of these paths, in this order (the comment above it is the authoritative version): -1. **Suspension** (worker contexts). The Python call raises +1. **Suspension** (eval requests on worker contexts). The Python call raises `SuspensionRequired`; the context thread returns `{suspended, CallbackId, State, {Name, Args}}` to `py_context`, which runs the registered fun (`execute/2` in `py_callback`), possibly serving nested calls meanwhile (`wait_for_callback/2`), and resumes with - the `resume_callback` NIF. -2. **Blocking callback pipe** (owngil contexts). The context thread writes + the `resume_callback` NIF, replaying the expression in the request's + namespace with the callback results cached. +2. **Inline** (call requests on worker contexts). The context thread sends + `{py_callback, CallbackId, Name, Args}` to `py_context` and keeps serving + its own request queue while it waits (`ctx_call_erlang_inline`), so the + fun's nested `py:call` runs on the same thread; the fun's result comes back + through the `context_callback_reply` NIF. Nothing is replayed. The fun runs + in a process bound to the context that shares the caller's local env. +3. **Blocking callback pipe** (owngil contexts). The context thread writes a request on a pipe and blocks; the `py_context` process has a dedicated handler (`callback_handler_loop/1`) that runs the fun and writes the response frame back with `context_write_callback_response`. -3. **Thread worker** (`c_src/py_thread_worker.c`): any Python thread that is +4. **Thread worker** (`c_src/py_thread_worker.c`): any Python thread that is not a context thread (`threading.Thread`, executors) asks the `py_thread_handler` coordinator for a handler process and talks to it over a pipe. There is also an async variant (`erlang.async_call`) using a diff --git a/docs/asyncio.md b/docs/asyncio.md index b4f2990..442e964 100644 --- a/docs/asyncio.md +++ b/docs/asyncio.md @@ -712,10 +712,9 @@ def sync_handler(): | Context | Mechanism | What blocks | |---------|-----------|-------------| | Async (`await erlang.sleep()`) | `asyncio.sleep()` via Erlang `send_after` | Yields to the event loop. The worker pthread is free to handle other tasks. | -| Sync from `py:exec` / `py:eval` | `erlang.call('_py_sleep', secs)` triggers suspension; the dirty scheduler is released and an Erlang `receive ... after` parks the caller | Caller's Erlang process. Dirty scheduler free for other work. | -| Sync from `py:call` (worker mode) | Falls back to `time.sleep`; replaying the Python frame around a suspension would change time-measurement semantics | The context's worker pthread for the sleep duration. Async NIF dispatch returns immediately so the BEAM dirty scheduler is **not** held; other Erlang processes and other contexts run normally. | +| Sync (`py:exec`, `py:eval`, `py:call`) | `erlang._call_blocking('_py_sleep', secs)`: the Erlang handler parks in `receive ... after`; the Python frame is not suspended, so nothing is replayed around the sleep | The context's worker pthread for the sleep duration. Async NIF dispatch returned immediately so the BEAM dirty scheduler is **not** held; other Erlang processes and other contexts run normally. | -In every case the BEAM dirty scheduler is freed during the sleep — the difference is which thread blocks (Erlang process, dirty scheduler, or worker pthread). +In every case the BEAM dirty scheduler is freed during the sleep — the difference is which thread blocks (Erlang process or worker pthread). #### asyncio.sleep(delay) @@ -996,7 +995,13 @@ def handler(): **Behavior:** - Blocks the current Python execution until the Erlang callback completes -- Code executes exactly once (no replay) +- From `py:call` (and from any Python thread) the code executes exactly once: + the context thread waits for the callback and serves its own context's + requests meanwhile, so a callback that calls `py:call` again runs on the + same context, at any nesting depth +- From `py:eval` the request is suspended and the expression is replayed + when the callback returns, with earlier callback results taken from a + cache: keep side effects out of code that runs before an `erlang.call` - The callback can release the dirty scheduler by using Erlang's `receive` (e.g., `erlang.sleep()`, `channel.receive()`) - Quick callbacks hold the dirty scheduler; callbacks that wait via `receive` release it diff --git a/priv/_erlang_impl/__init__.py b/priv/_erlang_impl/__init__.py index 0b69096..5bc3c69 100644 --- a/priv/_erlang_impl/__init__.py +++ b/priv/_erlang_impl/__init__.py @@ -272,18 +272,15 @@ def sleep(seconds): - Async (``await erlang.sleep()``) uses ``asyncio.sleep()``, which routes through Erlang's ``send_after`` timer. The coroutine yields to the event loop; the worker pthread handles other tasks. - - Sync from ``py:exec`` / ``py:eval`` calls - ``erlang.call('_py_sleep', seconds)``. The suspension machinery - releases the dirty scheduler and parks the caller's Erlang - process in a ``receive ... after``. - - Sync from ``py:call`` falls back to ``time.sleep`` — the worker - pthread blocks for the sleep duration. The BEAM dirty scheduler - is *not* held here either: the NIF dispatch returned immediately - and the caller is waiting in an Erlang ``receive`` on the - context process. Other Erlang processes and other contexts run - normally during the sleep. (Replaying a suspended Python frame - around ``time.time()`` would change time-measurement semantics, - which is why ``py:call`` doesn't take the suspension path.) + - Sync (``py:exec``, ``py:eval``, ``py:call``) calls + ``erlang._call_blocking('_py_sleep', seconds)``: the Erlang + handler parks in a ``receive ... after`` while the context's + worker pthread waits on the thread worker pipe. The Python frame + is not suspended, so nothing is replayed around the sleep. The + BEAM dirty scheduler is not held: the NIF dispatch returned + immediately and the caller is waiting in an Erlang ``receive`` + on the context process, so other Erlang processes and other + contexts run normally during the sleep. Args: seconds: Duration to sleep in seconds (float or int). @@ -306,20 +303,20 @@ def handler(): # Async context - return awaitable that uses Erlang timers return asyncio.sleep(seconds) except RuntimeError: - # Sync context - use erlang.call to truly suspend and free dirty scheduler + # Sync context: block on the thread worker path. A plain erlang.call + # would suspend the request and replay the Python frame around the + # sleep on resume, which changes time-measurement semantics. + import erlang + blocking = getattr(erlang, '_call_blocking', None) + if blocking is not None: + blocking('_py_sleep', seconds) + return try: - import erlang erlang.call('_py_sleep', seconds) - except BaseException as e: - # SuspensionRequiredException inherits from BaseException (not Exception). - # When suspension is triggered, the NIF would replay the entire Python - # function from the beginning after the callback completes. This causes - # issues with time measurement since time.time() is called again during - # replay. For sync sleep, we fall back to time.sleep() which blocks - # correctly from the caller's perspective. - # Note: This means the dirty scheduler is NOT freed during sync sleep - # when running in context_call mode. For proper dirty scheduler release - # in sync contexts, use py:exec/py:eval instead of py:call. + except BaseException: + # No blocking path (an erlang module without _call_blocking): + # SuspensionRequired inherits from BaseException; sleep here + # instead of letting the frame be replayed. time.sleep(seconds) diff --git a/priv/_erlang_impl/_loop.py b/priv/_erlang_impl/_loop.py index e9f0b23..1797353 100644 --- a/priv/_erlang_impl/_loop.py +++ b/priv/_erlang_impl/_loop.py @@ -79,7 +79,7 @@ class ErlangEventLoop(asyncio.AbstractEventLoop): # Use __slots__ for faster attribute access and reduced memory __slots__ = ( - '_pel', '_loop_capsule', '_uses_global_capsule', + '_pel', '_loop_capsule', '_uses_global_capsule', '_owns_capsule', '_readers', '_writers', '_callbacks_by_cid', # callback_id -> (callback, args, event_type) for O(1) dispatch '_fd_resources', # fd -> fd_key (shared fd_resource_t per fd) @@ -96,15 +96,17 @@ class ErlangEventLoop(asyncio.AbstractEventLoop): '_wake_pending', # coalesced wakeup flag for call_soon_threadsafe ) - def __init__(self): + def __init__(self, capsule=None): """Initialize the Erlang event loop. The event loop is backed by Erlang's scheduler via the py_event_loop C module. This provides direct access to the event loop without going through Erlang callbacks. - Each loop instance has its own isolated capsule for proper timer - and FD event routing. + Without ``capsule`` the loop attaches to the interpreter's default + (global) loop. The event loop pool passes a capsule over one of its + own loops so timers, fd events and the pending queue stay on that + loop; such a capsule is owned by Erlang and never destroyed here. """ # Detect execution mode for proper behavior self._execution_mode = detect_mode() @@ -130,7 +132,10 @@ def __init__(self): # Without this, Python-created loops would have their own pending queues # that never get processed because the worker doesn't know about them. self._uses_global_capsule = False - if hasattr(self._pel, '_get_global_loop_capsule'): + self._owns_capsule = False + if capsule is not None: + self._loop_capsule = capsule + elif hasattr(self._pel, '_get_global_loop_capsule'): try: self._loop_capsule = self._pel._get_global_loop_capsule() self._uses_global_capsule = True @@ -148,8 +153,10 @@ def __init__(self): # Fall back to creating a new loop if global not available self._loop_capsule = self._pel._loop_new() self._uses_global_capsule = False + self._owns_capsule = True else: self._loop_capsule = self._pel._loop_new() + self._owns_capsule = True # Store reference to this Python loop in the C struct # This enables process_ready_tasks to access the loop directly @@ -159,7 +166,7 @@ def __init__(self): # Also set reference on the global interpreter loop # This is needed for py_nif:submit_task which uses the global loop - if hasattr(self._pel, '_set_global_loop_ref'): + if self._uses_global_capsule and hasattr(self._pel, '_set_global_loop_ref'): try: self._pel._set_global_loop_ref(self) except RuntimeError: @@ -367,8 +374,9 @@ def close(self): except Exception: pass - # Destroy loop capsule (but not if using shared global capsule) - if not self._uses_global_capsule: + # Destroy loop capsule (not a shared one: the global capsule or a + # pool loop's, both owned by Erlang) + if getattr(self, '_owns_capsule', not self._uses_global_capsule): try: self._pel._loop_destroy(self._loop_capsule) except Exception: diff --git a/priv/_erlang_impl/_shm.py b/priv/_erlang_impl/_shm.py index 5b490a9..4711f59 100644 --- a/priv/_erlang_impl/_shm.py +++ b/priv/_erlang_impl/_shm.py @@ -27,6 +27,10 @@ erlang.call('_py_buffer_wait', id, read_pos) -> (write_pos, closed) erlang.call('_py_buffer_consumed', id, n) -> ok erlang.call('_py_buffer_state', id) -> (write_pos, closed) + +In the embedded interpreter these use erlang._call_blocking: a plain +erlang.call from a context request suspends the request and replays the +Python frame on resume, which would take from the ring twice. """ import collections @@ -57,6 +61,13 @@ def _erlang(): return erlang +def _call(name, *args): + """Call an Erlang flow-control callback without suspending the caller.""" + erlang = _erlang() + call = getattr(erlang, '_call_blocking', None) or erlang.call + return call(name, *args) + + def _atom(name): return _erlang().atom(name) @@ -168,13 +179,13 @@ def _refresh_header(self): def _wait_for_data(self): """Block until write position passes our read position or EOF.""" - wpos, closed = _erlang().call('_py_buffer_wait', self.id, self._rpos) + wpos, closed = _call('_py_buffer_wait', self.id, self._rpos) self._wpos = wpos self._closed = bool(closed) def _consumed(self, n): if n: - _erlang().call('_py_buffer_consumed', self.id, n) + _call('_py_buffer_consumed', self.id, n) def _available(self): return self._wpos - self._rpos diff --git a/src/py.erl b/src/py.erl index f2c682e..c564138 100644 --- a/src/py.erl +++ b/src/py.erl @@ -85,6 +85,7 @@ activate_venv/1, %% Process-local Python environment get_local_env/1, + put_local_env/2, deactivate_venv/0, venv_info/0, %% Execution info @@ -183,6 +184,20 @@ get_local_env(Ctx) when is_pid(Ctx) -> Ref end. +%% @private Seed the calling process with an existing local env. +%% +%% Used by a context when it spawns the process that runs an Erlang +%% callback, so the callback shares the suspended caller's Python +%% namespace instead of starting from an empty one. +-spec put_local_env(non_neg_integer(), reference()) -> ok. +put_local_env(InterpId, EnvRef) -> + Envs = case get(?LOCAL_ENV_KEY) of + M when is_map(M) -> M; + _ -> #{} + end, + put(?LOCAL_ENV_KEY, Envs#{InterpId => EnvRef}), + ok. + %%% ============================================================================ %%% Synchronous API %%% ============================================================================ diff --git a/src/py_context_embedded.erl b/src/py_context_embedded.erl index fddeb48..dd73771 100644 --- a/src/py_context_embedded.erl +++ b/src/py_context_embedded.erl @@ -576,7 +576,7 @@ handle_call_with_suspension(Ref, Module, Func, Args, Kwargs) -> case py_nif:context_call_async(Ref, self(), RequestId, Module, Func, Args, Kwargs) of {enqueued, RequestId} -> %% Async dispatch succeeded - wait for result message - wait_for_async_result(Ref, RequestId); + wait_for_async_result(Ref, RequestId, undefined); {error, Reason} -> {error, Reason} end. @@ -589,7 +589,7 @@ handle_eval_with_suspension(Ref, Code, Locals) -> case py_nif:context_eval_async(Ref, self(), RequestId, Code, Locals) of {enqueued, RequestId} -> %% Async dispatch succeeded - wait for result message - wait_for_async_result(Ref, RequestId); + wait_for_async_result(Ref, RequestId, undefined); {error, Reason} -> {error, Reason} end. @@ -600,7 +600,7 @@ handle_exec_with_async(Ref, Code) -> RequestId = make_ref(), case py_nif:context_exec_async(Ref, self(), RequestId, Code) of {enqueued, RequestId} -> - wait_for_async_result(Ref, RequestId); + wait_for_async_result(Ref, RequestId, undefined); {error, Reason} -> {error, Reason} end. @@ -621,15 +621,37 @@ handle_exec_with_async(Ref, Code) -> %% async results and only one wait_for_async_result/2 is in flight at %% a time, so the drain cannot consume the result of a concurrent live %% request. -wait_for_async_result(Ref, RequestId) -> +%% +%% EnvRef is the process-local env of the request (undefined when it has +%% none); a callback spawned for it inherits it. +%% +%% A call request that reaches erlang.call sends {py_callback, ...} and +%% waits inline, serving this context's queue meanwhile: the callback runs +%% here as for a suspension (nested requests are served through +%% wait_for_callback/2) and its result goes back with +%% context_callback_reply/3. +wait_for_async_result(Ref, RequestId, EnvRef) -> drain_stale_async_results(RequestId), receive {py_result, RequestId, Result} -> - process_async_result(Ref, Result) + process_async_result(Ref, Result, EnvRef); + {py_callback, CallbackId, FuncName, CallbackArgs} -> + CallbackResult = handle_callback_with_nested_receive(Ref, FuncName, CallbackArgs, EnvRef), + reply_inline_callback(Ref, CallbackId, CallbackResult), + wait_for_async_result(Ref, RequestId, EnvRef) after 300000 -> %% 5 minute timeout {error, async_timeout} end. +%% @private +reply_inline_callback(Ref, CallbackId, {ok, ResultBin}) -> + _ = py_nif:context_callback_reply(Ref, CallbackId, ResultBin), + ok; +reply_inline_callback(Ref, CallbackId, {error, Reason}) -> + ErrMsg = iolist_to_binary(io_lib:format("~p", [Reason])), + _ = py_nif:context_callback_reply(Ref, CallbackId, <<1, ErrMsg/binary>>), + ok. + %% @private drain_stale_async_results(CurrentId) -> receive @@ -642,12 +664,12 @@ drain_stale_async_results(CurrentId) -> %% @private %% Process the result from async dispatch %% Handles suspension, schedule markers, and normal results. -process_async_result(Ref, {suspended, _CallbackId, StateRef, {FuncName, CallbackArgs}}) -> - CallbackResult = handle_callback_with_nested_receive(Ref, FuncName, CallbackArgs), - resume_and_continue(Ref, StateRef, CallbackResult); -process_async_result(Ref, {schedule, CallbackName, CallbackArgs}) -> +process_async_result(Ref, {suspended, _CallbackId, StateRef, {FuncName, CallbackArgs}}, EnvRef) -> + CallbackResult = handle_callback_with_nested_receive(Ref, FuncName, CallbackArgs, EnvRef), + resume_and_continue(Ref, StateRef, CallbackResult, EnvRef); +process_async_result(Ref, {schedule, CallbackName, CallbackArgs}, _EnvRef) -> handle_schedule(Ref, CallbackName, CallbackArgs); -process_async_result(_Ref, Result) -> +process_async_result(_Ref, Result, _EnvRef) -> Result. %% @private @@ -658,7 +680,7 @@ handle_call_with_suspension_and_env(Ref, Module, Func, Args, Kwargs, EnvRef) -> Module, Func, Args, Kwargs, EnvRef) of {enqueued, RequestId} -> - wait_for_async_result(Ref, RequestId); + wait_for_async_result(Ref, RequestId, EnvRef); {error, Reason} -> {error, Reason} end. @@ -671,7 +693,7 @@ handle_eval_with_suspension_and_env(Ref, Code, Locals, EnvRef) -> case py_nif:context_eval_with_env_async(Ref, self(), RequestId, Code, Locals, EnvRef) of {enqueued, RequestId} -> - wait_for_async_result(Ref, RequestId); + wait_for_async_result(Ref, RequestId, EnvRef); {error, Reason} -> {error, Reason} end. @@ -685,7 +707,7 @@ handle_exec_with_async_and_env(Ref, Code, EnvRef) -> case py_nif:context_exec_with_env_async(Ref, self(), RequestId, Code, EnvRef) of {enqueued, RequestId} -> - wait_for_async_result(Ref, RequestId); + wait_for_async_result(Ref, RequestId, EnvRef); {error, Reason} -> {error, Reason} end. @@ -732,9 +754,21 @@ handle_schedule(_Ref, CallbackName, CallbackArgs) when is_binary(CallbackName) - %% Handle callback, allowing nested py:eval/call to be processed. %% We spawn a process to execute the callback so we can stay in a receive loop %% for nested calls while the callback runs. -handle_callback_with_nested_receive(Ref, FuncName, CallbackArgs) -> +%% +%% The callback runs on behalf of the suspended caller: its py:call/py:eval +%% are bound to this context (served by wait_for_callback/2, so nesting is +%% deadlock free at any depth) and it gets the caller's Python namespace, +%% so `py:call('__main__', F, Args)' from the callback sees the F the caller +%% defined with py:exec. +handle_callback_with_nested_receive(Ref, FuncName, CallbackArgs, EnvRef) -> Parent = self(), + InterpId = py_nif:context_interp_id(Ref), CallbackPid = spawn_link(fun() -> + ok = py_context_router:bind_context(Parent), + case EnvRef of + undefined -> ok; + _ -> ok = py:put_local_env(InterpId, EnvRef) + end, Result = try ArgsList = tuple_to_list(CallbackArgs), case py_callback:execute(FuncName, ArgsList) of @@ -824,16 +858,16 @@ wait_for_callback(Ref, CallbackPid) -> %% @private %% Resume suspended state, handle additional suspensions (nested callbacks) -resume_and_continue(Ref, StateRef, {ok, ResultBin}) -> +resume_and_continue(Ref, StateRef, {ok, ResultBin}, EnvRef) -> case py_nif:context_resume(Ref, StateRef, ResultBin) of {suspended, _CallbackId2, StateRef2, {FuncName2, Args2}} -> %% Another callback during resume - recursive handling - CallbackResult2 = handle_callback_with_nested_receive(Ref, FuncName2, Args2), - resume_and_continue(Ref, StateRef2, CallbackResult2); + CallbackResult2 = handle_callback_with_nested_receive(Ref, FuncName2, Args2, EnvRef), + resume_and_continue(Ref, StateRef2, CallbackResult2, EnvRef); FinalResult -> FinalResult end; -resume_and_continue(Ref, StateRef, {error, _} = Err) -> +resume_and_continue(Ref, StateRef, {error, _} = Err, _EnvRef) -> _ = py_nif:context_cancel_resume(Ref, StateRef), Err. diff --git a/src/py_event_loop_pool.erl b/src/py_event_loop_pool.erl index 4d0d16b..bc42f29 100644 --- a/src/py_event_loop_pool.erl +++ b/src/py_event_loop_pool.erl @@ -549,6 +549,7 @@ terminate(_Reason, State) -> Loops -> lists:foreach(fun({LoopRef, WorkerPid}) -> try py_event_worker:stop(WorkerPid) catch _:_ -> ok end, + try py_nif:event_loop_release_python_loop(LoopRef) catch _:_ -> ok end, try py_nif:event_loop_destroy(LoopRef) catch _:_ -> ok end end, tuple_to_list(Loops)) end, diff --git a/src/py_event_worker.erl b/src/py_event_worker.erl index 20963d0..0f0ec19 100644 --- a/src/py_event_worker.erl +++ b/src/py_event_worker.erl @@ -21,7 +21,11 @@ -record(state, { worker_id :: binary(), loop_ref :: reference(), - timers = #{} :: #{reference() => {reference(), non_neg_integer()}}, + %% TimerRef => {ErlTimerRef, LoopRef, CallbackId}. The loop is the one + %% named in the start_timer message: a loop without a worker of its own + %% borrows this one, so an expiry is dispatched to the loop that set it. + timers = #{} :: #{reference() | non_neg_integer() => + {reference(), reference(), non_neg_integer()}}, stats = #{select_count => 0, timer_count => 0, dispatch_count => 0} :: map() }). @@ -62,38 +66,42 @@ handle_info({select, FdRes, _Ref, ready_output}, State) -> maybe_send_task_ready(), {noreply, State}; -handle_info({start_timer, _LoopRef, DelayMs, CallbackId, TimerRef}, State) -> +handle_info({start_timer, LoopRef, DelayMs, CallbackId, TimerRef}, State) -> #state{timers = Timers} = State, ErlTimerRef = erlang:send_after(DelayMs, self(), {timeout, TimerRef}), - NewTimers = maps:put(TimerRef, {ErlTimerRef, CallbackId}, Timers), + NewTimers = maps:put(TimerRef, {ErlTimerRef, LoopRef, CallbackId}, Timers), {noreply, State#state{timers = NewTimers}}; handle_info({start_timer, DelayMs, CallbackId, TimerRef}, State) -> - #state{timers = Timers} = State, + #state{loop_ref = LoopRef, timers = Timers} = State, ErlTimerRef = erlang:send_after(DelayMs, self(), {timeout, TimerRef}), - NewTimers = maps:put(TimerRef, {ErlTimerRef, CallbackId}, Timers), + NewTimers = maps:put(TimerRef, {ErlTimerRef, LoopRef, CallbackId}, Timers), {noreply, State#state{timers = NewTimers}}; handle_info({cancel_timer, TimerRef}, State) -> #state{timers = Timers} = State, case maps:get(TimerRef, Timers, undefined) of undefined -> {noreply, State}; - {ErlTimerRef, _CallbackId} -> + {ErlTimerRef, _LoopRef, _CallbackId} -> erlang:cancel_timer(ErlTimerRef), NewTimers = maps:remove(TimerRef, Timers), {noreply, State#state{timers = NewTimers}} end; handle_info({timeout, TimerRef}, State) -> - #state{loop_ref = LoopRef, timers = Timers} = State, + #state{loop_ref = OwnLoopRef, timers = Timers} = State, case maps:get(TimerRef, Timers, undefined) of undefined -> {noreply, State}; - {_ErlTimerRef, CallbackId} -> + {_ErlTimerRef, LoopRef, CallbackId} -> py_nif:dispatch_timer(LoopRef, CallbackId), NewTimers = maps:remove(TimerRef, Timers), - %% Trigger event processing after timer dispatch - %% This ensures _run_once is called to handle the timer callback - maybe_send_task_ready(), + %% Trigger event processing after timer dispatch so _run_once + %% handles the timer callback. A borrowed loop is not served by + %% the task_ready loop of this worker, so drive it directly. + case LoopRef of + OwnLoopRef -> maybe_send_task_ready(); + _ -> _ = py_nif:process_ready_tasks(LoopRef) + end, {noreply, State#state{timers = NewTimers}} end; @@ -116,7 +124,7 @@ handle_info(task_ready, #state{loop_ref = LoopRef} = State) -> handle_info(_Info, State) -> {noreply, State}. terminate(_Reason, #state{timers = Timers}) -> - maps:foreach(fun(_TimerRef, {ErlTimerRef, _CallbackId}) -> + maps:foreach(fun(_TimerRef, {ErlTimerRef, _LoopRef, _CallbackId}) -> erlang:cancel_timer(ErlTimerRef) end, Timers), ok. diff --git a/src/py_nif.erl b/src/py_nif.erl index f4dddc3..1069b47 100644 --- a/src/py_nif.erl +++ b/src/py_nif.erl @@ -68,6 +68,7 @@ set_event_loop_priv_dir/1, event_loop_new/0, event_loop_destroy/1, + event_loop_release_python_loop/1, event_loop_set_router/2, event_loop_set_worker/2, event_loop_set_id/2, @@ -157,6 +158,7 @@ context_get_callback_pipe/1, context_write_callback_response/2, context_resume/3, + context_callback_reply/3, context_cancel_resume/2, context_get_event_loop/1, %% py_ref API (Python object references with interp_id) @@ -546,6 +548,15 @@ event_loop_new() -> event_loop_destroy(_LoopRef) -> ?NIF_STUB. +%% @doc Drop the loop's reference to its Python `ErlangEventLoop'. +%% +%% The Python loop holds a capsule that keeps the loop resource alive, so +%% a loop that ran tasks must release it before `event_loop_destroy/1' or +%% neither side is ever freed. +-spec event_loop_release_python_loop(reference()) -> ok | {error, term()}. +event_loop_release_python_loop(_LoopRef) -> + ?NIF_STUB. + %% @doc Set the router process for an event loop (legacy). %% The router receives enif_select messages and timer events. -spec event_loop_set_router(reference(), pid()) -> ok | {error, term()}. @@ -1274,6 +1285,17 @@ context_write_callback_response(_ContextRef, _Data) -> context_resume(_ContextRef, _StateRef, _Result) -> ?NIF_STUB. +%% @doc Deliver the result of an inline callback to the context thread. +%% +%% The thread sent `{py_callback, CallbackId, Name, Args}' from an +%% `erlang.call' in a call request and waits for this reply while it +%% serves the context's queue. Result is the frame the callback decoder +%% reads: a status byte (2 ok, 1 error) then the payload. +-spec context_callback_reply(reference(), non_neg_integer(), binary()) -> + ok | {error, term()}. +context_callback_reply(_ContextRef, _CallbackId, _Result) -> + ?NIF_STUB. + %% @doc Cancel a suspended context resume (cleanup on error). %% %% Called when callback execution fails and resume won't be called. diff --git a/test/py_event_loop_pool_SUITE.erl b/test/py_event_loop_pool_SUITE.erl index 6ae4737..d1e39a8 100644 --- a/test/py_event_loop_pool_SUITE.erl +++ b/test/py_event_loop_pool_SUITE.erl @@ -29,6 +29,9 @@ %% Ordering test test_tasks_execute_in_order/1, + %% Timer test + test_concurrent_sleeps_complete/1, + %% Exec/Eval tests test_exec_basic/1, test_eval_basic/1, @@ -47,6 +50,7 @@ all() -> test_spawn_task, test_concurrent_tasks, test_tasks_execute_in_order, + test_concurrent_sleeps_complete, test_exec_basic, test_eval_basic, test_exec_eval_namespace, @@ -172,6 +176,46 @@ test_tasks_execute_in_order(_Config) -> [{ok, 1.0}, {ok, 2.0}, {ok, 3.0}, {ok, 4.0}, {ok, 5.0}] = Results, ok. +%% ============================================================================ +%% Timer Test +%% ============================================================================ + +%% @doc Many tasks sleeping at once, spread over the pool, must all wake up. +%% +%% Regression: every pool loop used to schedule its timers on the global +%% loop and poll the global pending queue, and per-loop callback ids +%% collided there, so with 24 concurrent 50 ms sleeps only one task ever +%% resumed. Sleeps of 0 or 1 ms bypass the timer path, so the sleep here is +%% long enough to go through erlang:send_after. +test_concurrent_sleeps_complete(_Config) -> + TestDir = filename:join(code:lib_dir(erlang_python), "test"), + ok = py:exec(iolist_to_binary(io_lib:format( + "import sys; sys.path.insert(0, '~s')", [TestDir]))), + N = 32, + SleepMs = 100, + Self = self(), + %% One process per task: get_loop/0 binds a loop per calling process + Pids = [spawn_link(fun() -> + {ok, Loop} = py_event_loop_pool:get_loop(), + Ref = make_ref(), + ok = py_nif:submit_task(Loop, self(), Ref, <<"py_test_pool_sleep">>, + <<"nap">>, [SleepMs], #{}), + R = receive {async_result, Ref, X} -> X after 5000 -> timeout end, + Self ! {done, self(), Loop, R} + end) || _ <- lists:seq(1, N)], + Results = [receive {done, Pid, Loop, R} -> {Loop, R} end || Pid <- Pids], + Timeouts = [x || {_, timeout} <- Results], + ct:log("results: ~p", [Results]), + [] = Timeouts, + [{ok, _} = R || {_, R} <- Results], + Loops = lists:usort([Loop || {Loop, _} <- Results]), + ct:log("tasks ran on ~p loops", [length(Loops)]), + case maps:get(num_loops, py_event_loop_pool:get_stats(), 1) of + 1 -> ok; + _ -> true = length(Loops) > 1 + end, + ok. + %% ============================================================================ %% Exec/Eval Tests %% ============================================================================ diff --git a/test/py_reentrant_SUITE.erl b/test/py_reentrant_SUITE.erl index 77f447a..47af7be 100644 --- a/test/py_reentrant_SUITE.erl +++ b/test/py_reentrant_SUITE.erl @@ -27,7 +27,11 @@ test_async_call/1, test_callback_name_registry/1, test_etf_decode_safe/1, - test_reentrant_resume_stress/1 + test_reentrant_resume_stress/1, + test_call_reentrant_depths/1, + test_call_reentrant_repeated/1, + test_call_sequential_callbacks/1, + test_readme_reentrant_example/1 ]). all() -> @@ -44,7 +48,11 @@ all() -> test_async_call, test_callback_name_registry, test_etf_decode_safe, - test_reentrant_resume_stress + test_reentrant_resume_stress, + test_call_reentrant_depths, + test_call_reentrant_repeated, + test_call_sequential_callbacks, + test_readme_reentrant_example ]. init_per_suite(Config) -> @@ -74,6 +82,94 @@ end_per_testcase(_TestCase, _Config) -> try py:unregister_function(etf_probe_ok) catch _:_ -> ok end, try py:unregister_function(etf_probe_novel) catch _:_ -> ok end, try py:unregister_function(rs_double) catch _:_ -> ok end, + try py:unregister_function(nest_step) catch _:_ -> ok end, + ok. + +%%% ============================================================================ +%%% py:call as the outer call +%%% +%%% py:call used to run without suspension enabled: erlang.call blocked the +%%% context thread on the thread worker pipe, the callback's nested py:call +%%% landed on a busy context and the chain hung until the timeout, ending +%%% with "callback synchronisation lost; retry". py:eval chains worked. +%%% ============================================================================ + +reentrant_test_dir() -> + TestDir = filename:join(code:lib_dir(erlang_python), "test"), + ok = py:exec(iolist_to_binary(io_lib:format( + "import sys; sys.path.insert(0, '~s')", [TestDir]))). + +register_nest_step() -> + py:register_function(nest_step, fun([N, D]) -> + {ok, R} = py:call(py_test_reentrant, down, [N + 1, D - 1], #{}, 10000), + R + end). + +%% @doc py:call -> erlang.call -> py:call, nested 1, 2, 3 and 5 deep. +test_call_reentrant_depths(_Config) -> + reentrant_test_dir(), + register_nest_step(), + lists:foreach(fun(Depth) -> + {ok, Depth} = py:call(py_test_reentrant, down, [0, Depth], #{}, 10000) + end, [1, 2, 3, 5]), + py:unregister_function(nest_step), + ok. + +%% @doc One depth, many times: the nested py:call must not depend on which +%% context the scheduler picks. +test_call_reentrant_repeated(_Config) -> + reentrant_test_dir(), + register_nest_step(), + lists:foreach(fun(I) -> + {ok, I} = py:call(py_test_reentrant, down, [I - 1, 1], #{}, 10000) + end, lists:seq(1, 20)), + py:unregister_function(nest_step), + ok. + +%% @doc Several sequential erlang.call in one function called with py:call, +%% the last callback calling Python again (the hooks pattern). The function +%% runs once; each call blocks inline while the context serves the nested +%% py:call. +test_call_sequential_callbacks(_Config) -> + reentrant_test_dir(), + py:register_function(add_ten, fun([X]) -> X + 10 end), + py:register_function(multiply_by_two, fun([X]) -> X * 2 end), + py:register_function(subtract_five, fun([X]) -> + {ok, R} = py:call(py_test_reentrant, minus, [X, 5], #{}, 10000), + R + end), + lists:foreach(fun(I) -> + %% ((I + 10) * 2) - 5 + Expected = ((I + 10) * 2) - 5, + {ok, Expected} = py:call(py_test_reentrant, chain, [I], #{}, 10000) + end, lists:seq(1, 5)), + ok. + +%% @doc The README "Reentrant Callbacks" example, as written there: the +%% callback's py:call('__main__', double, ...) must see the double the +%% caller defined with py:exec. +test_readme_reentrant_example(_Config) -> + %% Register an Erlang function that calls Python + py:register_function(double_via_python, fun([X]) -> + {ok, Result} = py:call('__main__', double, [X]), + Result + end), + + %% Define Python functions + ok = py:exec(<<" +def double(x): + return x * 2 + +def process(x): + from erlang import call + # This calls Erlang, which calls Python's double() + doubled = call('double_via_python', x) + return doubled + 1 +">>), + + %% Test the full round-trip + {ok, 21} = py:call('__main__', process, [10]), + %% 10 -> double_via_python -> double(10)=20 -> +1 = 21 ok. %%% ============================================================================ diff --git a/test/py_test_pool_sleep.py b/test/py_test_pool_sleep.py new file mode 100644 index 0000000..a7146c7 --- /dev/null +++ b/test/py_test_pool_sleep.py @@ -0,0 +1,15 @@ +"""Helper for py_event_loop_pool_SUITE: a coroutine that sleeps on the loop. + +Sleeps of a few milliseconds or more go through the Erlang timer path +(erlang:send_after via the loop's worker); the suite submits many of them +across the pool at once. +""" + +import asyncio +import threading + + +async def nap(ms): + """Sleep ``ms`` milliseconds, then return the thread that resumed us.""" + await asyncio.sleep(ms / 1000) + return threading.get_ident() diff --git a/test/py_test_reentrant.py b/test/py_test_reentrant.py new file mode 100644 index 0000000..76451d6 --- /dev/null +++ b/test/py_test_reentrant.py @@ -0,0 +1,24 @@ +"""Helpers for py_reentrant_SUITE: functions called with py:call that call +back into Erlang with erlang.call, where the Erlang side calls py:call again. +""" + +import erlang + + +def down(n, depth): + """Recurse through Erlang: nest_step does py:call(down, [n + 1, depth - 1]).""" + if depth <= 0: + return n + return erlang.call('nest_step', n, depth) + + +def chain(x): + """Three sequential erlang.call in one function.""" + step1 = erlang.call('add_ten', x) + step2 = erlang.call('multiply_by_two', step1) + return erlang.call('subtract_five', step2) + + +def minus(x, y): + """Plain helper the subtract_five callback calls with py:call.""" + return x - y