From f735fd39296f4b180c78d54b49ac642b136234c9 Mon Sep 17 00:00:00 2001 From: ausimian Date: Thu, 16 Apr 2026 23:46:04 +0930 Subject: [PATCH] Fix SIGABRT/SIGSEGV from concurrent mx::eval dispatch MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit MLX is not thread-safe (ml-explore/mlx#2133). The Metal CommandEncoder is shared state — concurrent mx::eval calls from different OS threads crash with "A command encoder is already encoding to this command buffer" (SIGABRT) or SIGSEGV. The M14 soak tests were the first to exercise concurrent eval from multiple dirty-CPU scheduler threads. Fix: - Add emily::safe_eval() with a static mutex serialising all mx::eval calls (c_src/emily/tensor.hpp) - Remove set_default_stream calls from with_stream/2 — the NIF mutated MLX thread-local state which is unreliable under BEAM process migration between OS threads - Harden resolve_stream(-1) to avoid reading thread-local default Also adds M15.5 milestone for building MLX from source once upstream gains native thread-safety (0.32+). --- PLAN.md | 56 ++++++++++++++++++++++++++++++++++++++++++ RELEASE.md | 10 ++++++++ c_src/emily/tensor.hpp | 33 ++++++++++++++++++++++--- c_src/emily_nif.cpp | 4 +-- c_src/stream.cpp | 7 ++++-- lib/emily/stream.ex | 18 ++++---------- 6 files changed, 108 insertions(+), 20 deletions(-) diff --git a/PLAN.md b/PLAN.md index 9059661..50383ec 100644 --- a/PLAN.md +++ b/PLAN.md @@ -674,6 +674,31 @@ silent about it and ships a tested pattern. **Exit:** both patterns documented; concurrency soak demonstrates the streamed path is stable. +**Post-M14 note — eval serialisation:** +MLX is not thread-safe (ml-explore/mlx#2133). The Metal +`CommandEncoder` is shared state — concurrent `mx::eval` calls from +different OS threads crash with `"A command encoder is already +encoding to this command buffer"` (SIGABRT) or SIGSEGV from corrupted +encoder state. The M14 soak tests (8 concurrent workers via +`Task.async_stream`) exposed this by being the first code path to call +`mx::eval` from multiple dirty-CPU scheduler threads simultaneously. + +Fixed by: +1. **`emily::safe_eval()`**: a mutex-serialised `mx::eval` wrapper in + `c_src/emily/tensor.hpp`. All eval callsites route through it. + Graph-building ops (regular scheduler) remain lock-free. +2. **Removed `set_default_stream` from `with_stream/2`**: the NIF + mutated MLX thread-local state on BEAM scheduler threads, which + is unreliable since BEAM processes migrate between OS threads. + The process-dictionary-based stream routing was already correct. +3. **Hardened `resolve_stream(-1)`**: the -1 fallback now reads the + device default directly instead of the (potentially corrupted) + thread-local default. + +The mutex serialises Metal dispatch at the cost of true concurrent +GPU execution. See M15.5 (MLX upgrade) for the plan to restore +concurrency via MLX's native thread-local `CommandEncoder` support. + ### M15 — Native linalg `lu`, `svd`, `qr`, `cholesky`, `triangular_solve`, `eigh`, @@ -698,6 +723,37 @@ BinaryBackend-slow. MLX exposes most natively under `mx::linalg::*`. **Exit:** all `mx::linalg::*`-backed callbacks pass property suite; remaining `via_binary` linalg paths documented with rationale. +### M15.5 — MLX upgrade (build from source) + +Emily pins MLX 0.25.1 via pre-built binaries from `cocoa-xu/mlx-build`. +MLX gained native thread-safety on `main` in April 2026 (thread-local +`CommandEncoder` ml-explore/mlx#3348, `ThreadLocalStream` C++ API +ml-explore/mlx#3405), but neither fix is in any release yet (latest: +0.31.1). Building from source unblocks true concurrent Metal dispatch +and removes the `safe_eval` mutex introduced in the post-M14 fix. + +- **Build MLX from source** (from `main` or 0.32+ when released) + instead of fetching the pre-built 0.25.1 tarball. Extend `mix.exs` + to support an `MLX_SOURCE` env var pointing at a local MLX build. +- **Audit API changes** between 0.25.1 and target version. Emily's + C++ surface is narrow (core ops, linalg, streams, eval, allocator), + but six months of MLX releases may rename or remove functions. +- **Adopt `ThreadLocalStream` C++ API**: each BEAM dirty-CPU thread + gets its own MLX stream automatically, enabling true per-thread + Metal command queues without the eval mutex. +- **Remove `emily::safe_eval` mutex** once native thread-safety is + validated — revert to direct `mx::eval` calls. + +**Testing**: +- Amplified stress test (16 workers, 100 iterations) passes without + mutex under the new MLX build. +- Full test suite 20x with zero crashes. +- Benchmark `to_binary` latency under concurrent load: confirm + throughput improves vs. the mutex-serialised path. + +**Exit:** concurrent soak tests pass without mutex; MLX build-from-source +documented; stress test confirms concurrent Metal dispatch is stable. + ### M16 — Mixed-precision training bf16 activations + f32 master weights + loss scaling is the standard diff --git a/RELEASE.md b/RELEASE.md index b58f1c1..86d043f 100644 --- a/RELEASE.md +++ b/RELEASE.md @@ -2,6 +2,16 @@ ## Fixed +- **Fix SIGABRT/SIGSEGV from concurrent `mx::eval` dispatch.** MLX's + Metal `CommandEncoder` is not thread-safe (ml-explore/mlx#2133). + Concurrent `mx::eval` calls from BEAM dirty-CPU scheduler threads + triggered `"A command encoder is already encoding"` assertions or + SIGSEGV from corrupted encoder state. Fixed by serialising all + `mx::eval` calls through `emily::safe_eval()` (mutex in + `c_src/emily/tensor.hpp`). Also removed `set_default_stream` calls + from `with_stream/2` — the NIF mutated MLX thread-local state which + is unreliable under BEAM process migration. Hardened + `resolve_stream(-1)` to avoid reading the thread-local default. - Relax MNIST convergence canary threshold from 97% to 96% to eliminate stochastic flaps (observed 96.99% on occasional runs). The test is a sanity gate, not a performance benchmark. diff --git a/c_src/emily/tensor.hpp b/c_src/emily/tensor.hpp index 65bf24b..90d7a22 100644 --- a/c_src/emily/tensor.hpp +++ b/c_src/emily/tensor.hpp @@ -15,6 +15,7 @@ #include #include +#include #include #include #include @@ -60,12 +61,38 @@ unwrap_all(const std::vector> &tensors) { } // Resolve a stream index from Elixir into an mx::Stream. -// -1 (the sentinel for "no explicit stream") falls through to the -// thread-local default — backwards-compatible with pre-M14 code paths. +// -1 (the sentinel for "no explicit stream") returns the default GPU +// stream (index 0). We intentionally avoid mx::default_stream() here +// because that reads a thread-local which can be corrupted on BEAM +// scheduler threads — the BEAM migrates processes between OS threads, +// so thread-local state is unreliable. inline mx::Stream resolve_stream(int64_t stream_index) { if (stream_index < 0) - return mx::default_stream(mx::default_device()); + return mx::default_stream(mx::Device(mx::Device::DeviceType::gpu)); return mx::get_stream(static_cast(stream_index)); } +// MLX is not thread-safe (ml-explore/mlx#2133). In particular, the +// Metal CommandEncoder is shared state — concurrent mx::eval calls +// from different OS threads crash with "A command encoder is already +// encoding to this command buffer". BEAM dirty-CPU schedulers are a +// thread pool, so concurrent to_binary / eval NIF calls race. +// +// Serialise all mx::eval calls behind a single mutex until MLX gains +// native thread-safety (expected 0.32+, see ml-explore/mlx#3348). +inline std::mutex &eval_mutex() { + static std::mutex m; + return m; +} + +inline void safe_eval(mx::array &a) { + std::lock_guard lock(eval_mutex()); + mx::eval(a); +} + +inline void safe_eval(std::initializer_list arrays) { + std::lock_guard lock(eval_mutex()); + mx::eval(arrays); +} + } // namespace emily diff --git a/c_src/emily_nif.cpp b/c_src/emily_nif.cpp index d434dda..2accb08 100644 --- a/c_src/emily_nif.cpp +++ b/c_src/emily_nif.cpp @@ -79,7 +79,7 @@ FINE_NIF(from_binary, 0); fine::Term to_binary(ErlNifEnv *env, fine::ResourcePtr tensor, int64_t s) { auto stream = emily::resolve_stream(s); auto materialized = mx::contiguous(tensor->array, false, stream); - mx::eval(materialized); + emily::safe_eval(materialized); // Defensive: mx::contiguous is supposed to give a row-contiguous // layout. If it ever doesn't, we'd be aliasing a strided buffer and @@ -113,7 +113,7 @@ FINE_NIF(dtype, 0); // eval/1 — force evaluation of the lazy graph rooted at this tensor. // Dirty CPU: waits for MLX to finish. fine::Ok<> eval(ErlNifEnv *, fine::ResourcePtr tensor) { - mx::eval(tensor->array); + emily::safe_eval(tensor->array); return fine::Ok<>{}; } FINE_NIF(eval, ERL_NIF_DIRTY_JOB_CPU_BOUND); diff --git a/c_src/stream.cpp b/c_src/stream.cpp index 83e89f7..9f67c46 100644 --- a/c_src/stream.cpp +++ b/c_src/stream.cpp @@ -31,8 +31,11 @@ int64_t new_stream(ErlNifEnv *, fine::Atom device_atom) { FINE_NIF(new_stream, 0); // set_default_stream/1 — set the thread-local default stream. -// Belt-and-suspenders for Emily.Stream.with_stream/2; every op NIF -// also receives the stream index explicitly. +// +// WARNING: BEAM processes migrate between OS threads, so thread-local +// state is unreliable. with_stream/2 no longer calls this — it routes +// streams via the process dictionary instead. This NIF is retained for +// advanced use cases but should be avoided in normal code. fine::Ok<> set_default_stream(ErlNifEnv *, int64_t stream_index) { mx::set_default_stream(mx::get_stream(static_cast(stream_index))); return fine::Ok<>{}; diff --git a/lib/emily/stream.ex b/lib/emily/stream.ex index 80d39b1..2e3c9db 100644 --- a/lib/emily/stream.ex +++ b/lib/emily/stream.ex @@ -23,10 +23,9 @@ defmodule Emily.Stream do `with_stream/2` stores the stream index in the process dictionary. `Emily.Backend` reads it via `Process.get(:emily_stream, -1)` and passes it as an explicit argument to every NIF call. Each NIF - resolves the index to an `mx::Stream` (or falls back to the - thread-local default when the index is -1). This avoids the - thread-local race that would occur if we relied solely on - `set_default_stream` — BEAM processes can migrate between OS + resolves the index to an `mx::Stream` (or falls back to the default + GPU stream when the index is -1). This is process-dictionary based, + not thread-local, because BEAM processes migrate between OS scheduler threads between NIF calls. ## Concurrent serving patterns @@ -73,18 +72,13 @@ defmodule Emily.Stream do Execute `fun` with the given stream as the default for MLX ops. Stores the stream index in the process dictionary so that - `Emily.Backend` passes it to every NIF call. Also sets the - thread-local default stream as a belt-and-suspenders measure for - code that calls `Emily.Native` directly. - - The previous stream (if any) is restored in an `after` block, so - nesting is safe. + `Emily.Backend` passes it to every NIF call. The previous stream + (if any) is restored in an `after` block, so nesting is safe. """ @spec with_stream(t(), (-> result)) :: result when result: var def with_stream(%__MODULE__{index: index}, fun) when is_function(fun, 0) do prev = Process.get(:emily_stream) Process.put(:emily_stream, index) - Emily.Native.set_default_stream(index) try do fun.() @@ -93,8 +87,6 @@ defmodule Emily.Stream do nil -> Process.delete(:emily_stream) idx -> Process.put(:emily_stream, idx) end - - if prev, do: Emily.Native.set_default_stream(prev) end end