From 39e710efb8d5f3886122fe084e9affbf35a9e674 Mon Sep 17 00:00:00 2001 From: Nathan Clack Date: Wed, 23 Sep 2026 17:04:02 -0700 Subject: [PATCH 1/9] Add native macOS CPU platform support --- .github/workflows/cpu.yml | 21 +- CMakeLists.txt | 15 +- README.md | 17 +- cmake/Helpers.cmake | 5 +- docs/pipeline.md | 19 +- docs/prefetch.md | 7 +- docs/troubleshooting.md | 6 +- src/CMakeLists.txt | 15 +- src/damacy.h | 10 +- src/platform/numa.mach.c | 61 +++ src/platform/platform.mach.c | 295 ++++++++++++ src/store/metadata_store_async.h | 10 +- src/store/metadata_store_async.mach.c | 653 ++++++++++++++++++++++++++ tests/CMakeLists.txt | 4 + tests/test_array_meta.c | 1 + tests/test_metadata_store_async.c | 140 ++++++ tests/test_platform_mach.c | 47 ++ tests/test_store_latency.c | 1 + 18 files changed, 1297 insertions(+), 30 deletions(-) create mode 100644 src/platform/numa.mach.c create mode 100644 src/platform/platform.mach.c create mode 100644 src/store/metadata_store_async.mach.c create mode 100644 tests/test_platform_mach.c diff --git a/.github/workflows/cpu.yml b/.github/workflows/cpu.yml index 20699e57..54e2cafe 100644 --- a/.github/workflows/cpu.yml +++ b/.github/workflows/cpu.yml @@ -8,18 +8,27 @@ on: jobs: cpu: - runs-on: ubuntu-24.04 + strategy: + fail-fast: false + matrix: + os: [ubuntu-24.04, macos-26] + runs-on: ${{ matrix.os }} timeout-minutes: 15 steps: - uses: actions/checkout@v4 - uses: actions/setup-python@v5 with: python-version: '3.11' - - name: Install dependencies + - name: Install Linux dependencies + if: runner.os == 'Linux' run: | sudo apt-get update sudo apt-get install -y cmake ninja-build pkg-config liburing-dev libzstd-dev libblosc-dev - python -m pip install uv pytest pytest-cov numpy + - name: Install macOS dependencies + if: runner.os == 'macOS' + run: brew install cmake ninja pkg-config zstd c-blosc + - name: Install Python dependencies + run: python -m pip install uv pytest pytest-cov numpy - name: Build without CUDA run: | cmake --preset cpu -DDAMACY_PYTHON=ON -DPython_EXECUTABLE="$(command -v python)" -DCMAKE_DISABLE_FIND_PACKAGE_CUDAToolkit=ON @@ -32,9 +41,13 @@ jobs: run: | python - <<'PY' import subprocess + import sys from damacy import _native assert not _native.CUDA_ENABLED - dependencies = subprocess.check_output(['ldd', _native.__file__], text=True).lower() + command = ['otool', '-L'] if sys.platform == 'darwin' else ['ldd'] + dependencies = subprocess.check_output([*command, _native.__file__], text=True).lower() + if sys.platform == 'darwin': + assert 'liburing' not in dependencies print(dependencies) assert all(name not in dependencies for name in ('libcuda', 'libcudart', 'libnvcomp')) PY diff --git a/CMakeLists.txt b/CMakeLists.txt index 667bd807..1705e211 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -15,14 +15,25 @@ include(Warnings) include(Helpers) include(Fuzz) # declares DAMACY_FUZZ and applies fuzz-mode flags -option(DAMACY_CUDA "Build the CUDA executor" ON) +# CUDA is unavailable on current macOS toolchains; keep Linux's default. +if(APPLE) + set(DAMACY_CUDA_DEFAULT OFF) +else() + set(DAMACY_CUDA_DEFAULT ON) +endif() +option(DAMACY_CUDA "Build the CUDA executor" ${DAMACY_CUDA_DEFAULT}) if(DAMACY_FUZZ) set(DAMACY_CUDA OFF) endif() +if(APPLE AND DAMACY_CUDA) + message(FATAL_ERROR "macOS supports the CPU executor; set DAMACY_CUDA=OFF") +endif() find_package(Threads REQUIRED) find_package(PkgConfig REQUIRED) -pkg_check_modules(LIBURING REQUIRED IMPORTED_TARGET liburing) +if(NOT APPLE) + pkg_check_modules(LIBURING REQUIRED IMPORTED_TARGET liburing) +endif() if(DAMACY_CUDA) if(NOT DEFINED CMAKE_CUDA_ARCHITECTURES) diff --git a/README.md b/README.md index bf9765b7..9b333f05 100644 --- a/README.md +++ b/README.md @@ -61,7 +61,8 @@ while allowing the pool to reuse its storage. See [Pipeline composition](docs/pipeline.md) for the complete component, limit, ownership, and build contracts. Build a CPU Python package from source -with the Linux dependencies `liburing`, `libzstd`, and `libblosc` installed: +with `libzstd` and `libblosc` installed (plus `liburing` on Linux). On macOS, +install dependencies with `brew install cmake ninja pkg-config zstd c-blosc`: ```sh pip install . --config-settings=cmake.define.DAMACY_CUDA=OFF @@ -161,23 +162,27 @@ If you have data that uses one of the unsupported codecs and you'd like it added ## Runtime dependencies -All builds require Linux async metadata I/O and CPU codec libraries. CUDA -builds additionally link the NVIDIA driver and nvCOMP. Build with +CPU builds support Linux and macOS and require the CPU codec libraries. +Linux uses io_uring for async metadata I/O; macOS uses a POSIX worker pool. +CUDA builds additionally link the NVIDIA driver and nvCOMP. Build with `DAMACY_CUDA=OFF` to import and run on a host without a CUDA driver. | Library | Used by | How it is loaded | |---|---|---| -| `liburing` | CPU and CUDA: async metadata I/O | normal dynamic loader | +| `liburing` | Linux async metadata I/O | normal dynamic loader | | `libzstd`, `libblosc` | CPU decoding, included in both builds | normal dynamic loader | | `libcuda.so.1`, nvCOMP | CUDA builds | driver loader; nvCOMP may be linked statically | | `libnuma.so.1` | Optional CUDA placement and host affinity | `dlopen`; absence disables placement | | `libcufile.so.0` | Optional CUDA GPUDirect Storage | `dlopen`; requires `DAMACY_ENABLE_GDS=ON` | | `libmount.so.1`, `libudev.so.1` | cuFile initialization when GDS is used | dynamic loader | -Metadata reads require a Linux kernel with the io_uring operations damacy uses: +On Linux, metadata reads require a kernel with the io_uring operations damacy uses: `STATX`, `OPENAT2`, `READ`, and `CLOSE`. If the kernel does not advertise those operations, pipeline construction fails instead of falling back to a thread -pool. +pool. On macOS, `metadata_io_concurrency` sets the number of metadata workers; +bulk reads use the shared POSIX file backend. NUMA placement and CPU affinity +are unavailable. CUDA defaults off on macOS and cannot be enabled there. +See [native build instructions](docs/pipeline.md#build-without-cuda). GDS notes: diff --git a/cmake/Helpers.cmake b/cmake/Helpers.cmake index d8b19731..9648cfa4 100644 --- a/cmake/Helpers.cmake +++ b/cmake/Helpers.cmake @@ -31,7 +31,8 @@ function(add_cuda_lib TARGET) add_src_lib(${TARGET} ${ARGN}) endfunction() -# Appends basename.win32.c or basename.posix.c (or .darwin.c if it exists) +# Appends basename.win32.c or basename.posix.c (.mach.c on macOS, when +# present, with a legacy .darwin.c fallback) # to TARGET's sources. BASENAME defaults to TARGET. function(add_platform_sources TARGET) if(ARGC GREATER 1) @@ -41,6 +42,8 @@ function(add_platform_sources TARGET) endif() if(WIN32) target_sources(${TARGET} PRIVATE ${BASE}.win32.c) + elseif(APPLE AND EXISTS "${CMAKE_CURRENT_SOURCE_DIR}/${BASE}.mach.c") + target_sources(${TARGET} PRIVATE ${BASE}.mach.c) elseif(APPLE AND EXISTS "${CMAKE_CURRENT_SOURCE_DIR}/${BASE}.darwin.c") target_sources(${TARGET} PRIVATE ${BASE}.darwin.c) else() diff --git a/docs/pipeline.md b/docs/pipeline.md index 20080674..5bc81c28 100644 --- a/docs/pipeline.md +++ b/docs/pipeline.md @@ -180,6 +180,14 @@ On Linux, install C/C++ build tools, CMake, Ninja, pkg-config, liburing, zstd, and C-Blosc development packages. For example, Ubuntu packages are `build-essential cmake ninja-build pkg-config python3-dev liburing-dev libzstd-dev libblosc-dev`. +On macOS, install the native dependencies with Homebrew: + +```sh +brew install cmake ninja pkg-config zstd c-blosc uv +``` + +Then build and test on either platform: + ```sh cmake --preset cpu cmake --build build @@ -196,10 +204,13 @@ pip install . --config-settings=cmake.define.DAMACY_CUDA=OFF ``` The extension has no CUDA or nvCOMP dependency in this configuration. -`CudaExecutor` reports that CUDA support was not built. CPU builds retain the -Linux io_uring metadata requirements. CUDA support is enabled by default; -turning it on builds both executors and requires the CUDA toolkit, nvCOMP, and -a runtime NVIDIA driver. GDS requires a CUDA build. +`CudaExecutor` reports that CUDA support was not built. Linux retains its +io_uring metadata requirements. macOS uses a POSIX metadata worker pool, +with one worker per `metadata_io_concurrency`, and shared POSIX bulk reads. +CMake selects `platform.mach.c` and `numa.mach.c`; NUMA placement and CPU +affinity are unavailable. CUDA defaults off on macOS and cannot be enabled. +On Linux CUDA defaults on, builds both executors, and requires the CUDA toolkit, +nvCOMP, and a runtime NVIDIA driver. GDS requires a CUDA build. ## Future queries diff --git a/docs/prefetch.md b/docs/prefetch.md index a6a5ddc0..f18f8966 100644 --- a/docs/prefetch.md +++ b/docs/prefetch.md @@ -155,7 +155,12 @@ sets the metadata request-concurrency budget; the ring allocates enough entries internally for multi-step requests and driver wakeups. At startup the driver requires kernel support for `IORING_OP_STATX`, `IORING_OP_OPENAT2`, `IORING_OP_READ`, and `IORING_OP_CLOSE`; there is no thread-pool fallback in -the current build. +the Linux build. + +On macOS, CMake selects a POSIX worker-pool backend instead of io_uring. +`metadata_io_concurrency` is the worker count; each worker completes one +stat/open/read/close request at a time. Shutdown drains accepted requests. +Operation histograms measure syscall duration on macOS, excluding queue wait. The default metadata concurrency is 32. Treat much deeper values as storage tuning: they are useful when metadata operations have real latency, but each diff --git a/docs/troubleshooting.md b/docs/troubleshooting.md index 2db80869..8cfc3563 100644 --- a/docs/troubleshooting.md +++ b/docs/troubleshooting.md @@ -17,12 +17,14 @@ the calling thread yet. Two fixes: ## Pipeline construction fails during metadata I/O setup -The CPU and CUDA metadata paths use io_uring for Zarr metadata and shard indexes. At construction damacy requires kernel support for +On Linux, the CPU and CUDA metadata paths use io_uring for Zarr metadata and shard indexes. At construction damacy requires kernel support for `IORING_OP_STATX`, `IORING_OP_OPENAT2`, `IORING_OP_READ`, and `IORING_OP_CLOSE`. If ring creation or the operation probe fails, pipeline construction fails without a thread-pool fallback. Check the native log for the exact io_uring failure. -On supported kernels, an unusually high `metadata_io_concurrency` can also +macOS uses POSIX metadata worker threads and has no io_uring dependency. +An unusually high `metadata_io_concurrency` can exhaust thread resources on +macOS. On either platform it can also stress process file-descriptor limits because each in-flight metadata read can hold an open fd. The default is 64; for much deeper settings, check `ulimit -n` and remember to multiply by ranks per node. diff --git a/src/CMakeLists.txt b/src/CMakeLists.txt index 31af977c..024d3e65 100644 --- a/src/CMakeLists.txt +++ b/src/CMakeLists.txt @@ -32,7 +32,8 @@ add_src_lib(lru SOURCES util/lru.c util/lru.h LINKS log) add_src_lib(pool SOURCES util/pool.c util/pool.h LINKS log platform) add_src_lib(path_intern SOURCES util/path_intern.c util/path_intern.h LINKS hash) -# Platform abstraction (posix-only for now; chucky-style platform-split lives here). +# Platform abstraction: CMake selects Mach implementations on macOS and +# POSIX implementations on Linux. Shared POSIX file I/O works on both. # platform.posix.c uses dlopen for platform_dl*; numa.posix.c uses libnuma via # platform_dl* so a single binary works on hosts with or without libnuma installed. add_src_lib(platform @@ -188,10 +189,18 @@ else() target_sources(store PRIVATE store/store_fs_gds.h store/store_fs_gds_stub.c) endif() +if(APPLE) + set(METADATA_STORE_IMPL store/metadata_store_async.mach.c) +else() + set(METADATA_STORE_IMPL store/metadata_store_async.c) +endif() add_src_lib(metadata_store_async - SOURCES store/metadata_store_async.h store/metadata_store_async.c - LINKS store log platform PkgConfig::LIBURING + SOURCES store/metadata_store_async.h ${METADATA_STORE_IMPL} + LINKS store log platform ) +if(NOT APPLE) + target_link_libraries(metadata_store_async PUBLIC PkgConfig::LIBURING) +endif() if(TARGET numa) target_link_libraries(metadata_store_async PUBLIC numa) target_compile_definitions( diff --git a/src/damacy.h b/src/damacy.h index 17bf9c15..f832fbb0 100644 --- a/src/damacy.h +++ b/src/damacy.h @@ -114,8 +114,8 @@ extern "C" uint32_t n_io_threads; // Metadata request concurrency for array metadata, shard indexes, and // chunk-layout probes. The Linux metadata path uses this as an io_uring - // request-depth budget, not as a host thread count. Required: must be > 0 - // and no larger than DAMACY_MAX_METADATA_IO_CONCURRENCY. + // request-depth budget; macOS uses this many metadata worker threads. + // Required: must be > 0 and no larger than DAMACY_MAX_METADATA_IO_CONCURRENCY. uint32_t metadata_io_concurrency; uint32_t n_array_meta_cache; @@ -341,9 +341,9 @@ extern "C" uint64_t read_active; uint64_t read_max_active; } metadata_backend; - // Real, measured submit->completion latency of the io_uring metadata ops, - // broken out by op kind (statx/open/read/close). Distinct from the injected - // synthetic latency in metadata_latency above. Buckets are log2-scale on ns + // Measured metadata-operation latency: submit-to-completion on Linux, + // syscall duration on macOS, by kind (stat/open/read/close). Distinct from + // the injected synthetic latency in metadata_latency above. Buckets are log2-scale on ns // (bucket i: floor(log2(ns)) == i); percentiles are derived from them. struct { diff --git a/src/platform/numa.mach.c b/src/platform/numa.mach.c new file mode 100644 index 00000000..33b8811f --- /dev/null +++ b/src/platform/numa.mach.c @@ -0,0 +1,61 @@ +// macOS does not expose Linux NUMA-node or CPU-mask affinity controls. +#include "platform/numa.h" +#include +#include + +int +platform_numa_available(void) +{ + return 0; +} +int +platform_numa_max_node(void) +{ + return -1; +} +int +platform_numa_node_cpu_mask(int node, struct platform_cpu_mask* out) +{ + (void)node; + if (out) + memset(out, 0, sizeof(*out)); + return 1; +} +int +platform_thread_affinity_get(struct platform_cpu_mask* out) +{ + if (out) + memset(out, 0, sizeof(*out)); + return 1; +} +int +platform_thread_affinity_set(const struct platform_cpu_mask* mask) +{ + (void)mask; + return 1; +} +int +platform_cpu_mask_is_empty(const struct platform_cpu_mask* mask) +{ + if (!mask) + return 1; + for (size_t i = 0; i < sizeof(mask->bytes); ++i) + if (mask->bytes[i]) + return 0; + return 1; +} +int +platform_cpu_mask_describe(const struct platform_cpu_mask* mask, + int* first, + int* last, + int* count) +{ + (void)mask; + if (first) + *first = -1; + if (last) + *last = -1; + if (count) + *count = 0; + return 1; +} diff --git a/src/platform/platform.mach.c b/src/platform/platform.mach.c new file mode 100644 index 00000000..15a08aae --- /dev/null +++ b/src/platform/platform.mach.c @@ -0,0 +1,295 @@ +#include "platform/platform.h" + +#include +#include +#include +#include +#include +#include +#include +#include +#include + +size_t +platform_page_size(void) +{ + long size = sysconf(_SC_PAGESIZE); + return size > 0 ? (size_t)size : 0; +} + +size_t +platform_page_alignment(void) +{ + size_t ps = platform_page_size(); + return ps > 0 ? ps : 4096; +} + +size_t +platform_available_memory(void) +{ + vm_statistics64_data_t vm; + mach_msg_type_number_t count = HOST_VM_INFO64_COUNT; + mach_port_t host = mach_host_self(); + kern_return_t rc = + host_statistics64(host, HOST_VM_INFO64, (host_info64_t)&vm, &count); + mach_port_deallocate(mach_task_self(), host); + if (rc != KERN_SUCCESS) + return 0; + return ((size_t)vm.free_count + vm.inactive_count) * platform_page_size(); +} + +void* +platform_aligned_alloc(size_t alignment, size_t size) +{ + return aligned_alloc(alignment, size); +} + +void +platform_aligned_free(void* ptr) +{ + free(ptr); +} + +void +platform_sleep_ns(int64_t ns) +{ + struct timespec ts = { + .tv_sec = ns / 1000000000LL, + .tv_nsec = ns % 1000000000LL, + }; + nanosleep(&ts, NULL); +} + +static int64_t +monotonic_ns(void) +{ + struct timespec now; + clock_gettime(CLOCK_MONOTONIC, &now); + return (int64_t)now.tv_sec * 1000000000LL + now.tv_nsec; +} + +float +platform_toc(struct platform_clock* clock) +{ + int64_t now = monotonic_ns(); + float elapsed = (float)((now - clock->last_ns) / 1e9); + clock->last_ns = now; + return elapsed; +} + +void +platform_localtime(const time_t* t, struct tm* out) +{ + localtime_r(t, out); +} + +void +platform_call_once(platform_once* flag, void (*fn)(void)) +{ + pthread_once(flag, fn); +} + +struct platform_thread +{ + pthread_t handle; + void (*fn)(void*); + void* arg; +}; + +static void* +thread_trampoline(void* p) +{ + struct platform_thread* t = (struct platform_thread*)p; + t->fn(t->arg); + return NULL; +} + +struct platform_thread* +platform_thread_start(void (*fn)(void*), void* arg) +{ + struct platform_thread* t = + (struct platform_thread*)calloc(1, sizeof(struct platform_thread)); + if (!t) + return NULL; + t->fn = fn; + t->arg = arg; + if (pthread_create(&t->handle, NULL, thread_trampoline, t) != 0) { + free(t); + return NULL; + } + return t; +} + +int +platform_thread_join(struct platform_thread* t) +{ + if (!t) + return -1; + int rc = pthread_join(t->handle, NULL); + free(t); + return rc; +} + +struct platform_mutex +{ + pthread_mutex_t m; +}; + +struct platform_mutex* +platform_mutex_new(void) +{ + struct platform_mutex* m = + (struct platform_mutex*)calloc(1, sizeof(struct platform_mutex)); + if (!m) + return NULL; + if (pthread_mutex_init(&m->m, NULL) != 0) { + free(m); + return NULL; + } + return m; +} + +void +platform_mutex_free(struct platform_mutex* m) +{ + if (!m) + return; + pthread_mutex_destroy(&m->m); + free(m); +} + +void +platform_mutex_lock(struct platform_mutex* m) +{ + pthread_mutex_lock(&m->m); +} + +void +platform_mutex_unlock(struct platform_mutex* m) +{ + pthread_mutex_unlock(&m->m); +} + +struct platform_cond +{ + pthread_cond_t c; +}; + +struct platform_cond* +platform_cond_new(void) +{ + struct platform_cond* c = + (struct platform_cond*)calloc(1, sizeof(struct platform_cond)); + if (!c) + return NULL; + if (pthread_cond_init(&c->c, NULL) != 0) { + free(c); + return NULL; + } + return c; +} + +void +platform_cond_free(struct platform_cond* c) +{ + if (!c) + return; + pthread_cond_destroy(&c->c); + free(c); +} + +void +platform_cond_wait(struct platform_cond* c, struct platform_mutex* m) +{ + pthread_cond_wait(&c->c, &m->m); +} + +int +platform_cond_timedwait_ms(struct platform_cond* c, + struct platform_mutex* m, + int timeout_ms) +{ + struct timespec ts; + clock_gettime(CLOCK_REALTIME, &ts); + ts.tv_sec += timeout_ms / 1000; + ts.tv_nsec += (long)(timeout_ms % 1000) * 1000000L; + if (ts.tv_nsec >= 1000000000L) { + ts.tv_sec += 1; + ts.tv_nsec -= 1000000000L; + } + return pthread_cond_timedwait(&c->c, &m->m, &ts) == ETIMEDOUT ? 1 : 0; +} + +void +platform_cond_broadcast(struct platform_cond* c) +{ + pthread_cond_broadcast(&c->c); +} + +void +platform_cpu_pause(void) +{ +#if defined(__i386__) || defined(__x86_64__) + __asm__ __volatile__("pause" ::: "memory"); +#elif defined(__aarch64__) || defined(__arm__) + __asm__ __volatile__("yield" ::: "memory"); +#else + sched_yield(); +#endif +} + +void +platform_yield(void) +{ + sched_yield(); +} + +int +platform_default_thread_count(void) +{ + long n = sysconf(_SC_NPROCESSORS_ONLN); + return n > 0 ? (int)n : 1; +} + +void* +platform_dlopen(const char* path) +{ + return dlopen(path, RTLD_LAZY | RTLD_LOCAL); +} + +void* +platform_dlsym(void* handle, const char* name) +{ + return dlsym(handle, name); +} + +void +platform_dlclose(void* handle) +{ + if (!handle) + return; + dlclose(handle); +} + +const char* +platform_dlerror(void) +{ + return dlerror(); +} + +const char* +platform_getenv(const char* name) +{ + return getenv(name); +} + +uint64_t +platform_max_open_files(void) +{ + struct rlimit r; + if (getrlimit(RLIMIT_NOFILE, &r) != 0) + return 0; + if (r.rlim_cur == RLIM_INFINITY) + return 0; + return (uint64_t)r.rlim_cur; +} diff --git a/src/store/metadata_store_async.h b/src/store/metadata_store_async.h index 00a91e8f..2ce3816a 100644 --- a/src/store/metadata_store_async.h +++ b/src/store/metadata_store_async.h @@ -14,6 +14,11 @@ extern "C" struct metadata_store_async; struct numa_resolved; + // Successful submissions invoke the callback exactly once, on a backend + // thread. Callbacks may overlap and may submit more work. The read callback + // owns data and must free it (NULL for a successful empty read). + // Callers must stop submitting before destroy; destroy drains all accepted + // requests and joins callbacks. Do not destroy from a callback. typedef void (*metadata_store_read_cb)(void* user, enum damacy_status status, void* data, @@ -48,7 +53,8 @@ extern "C" uint64_t read_max_active; }; -// Log2-scale histogram of measured submit->completion latency. Bucket i holds +// Log2 histogram: submit-to-completion on Linux, syscall duration on macOS. +// Bucket i holds // ops whose latency in ns has floor(log2(ns)) == i (bucket 0 also catches 0 ns, // the last bucket catches everything above). Percentiles are derived from the // raw buckets at report time so the estimator can change without an ABI change. @@ -63,7 +69,7 @@ extern "C" uint64_t buckets[METADATA_OP_LATENCY_NBUCKETS]; }; - // Indexed by enum op_kind: statx, open, read, close. + // Indexed by enum op_kind: stat/statx, open, read, close. struct metadata_store_async_op_latency_stats { struct metadata_store_async_op_latency_kind diff --git a/src/store/metadata_store_async.mach.c b/src/store/metadata_store_async.mach.c new file mode 100644 index 00000000..3c0f7767 --- /dev/null +++ b/src/store/metadata_store_async.mach.c @@ -0,0 +1,653 @@ +// macOS metadata backend: bounded POSIX workers behind the existing async API. +// Bulk data still uses the unchanged store_fs/io_queue path. +#include "store/metadata_store_async.h" +#include "log/log.h" +#include "platform/platform.h" +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +enum request_kind +{ + REQ_READ_FILE, + REQ_READ, + REQ_STAT, +}; + +enum op_kind +{ + OP_STATX, + OP_OPEN, + OP_READ, + OP_CLOSE, +}; + +enum latency_op_kind +{ + LATENCY_OP_STAT, + LATENCY_OP_SUBMIT, +}; + +struct metadata_metrics +{ + struct + { + _Atomic uint64_t ops; + _Atomic uint64_t stat_ops; + _Atomic uint64_t submit_ops; + _Atomic uint64_t active; + _Atomic uint64_t max_active; + _Atomic uint64_t total_sleep_ns; + _Atomic uint64_t max_sleep_ns; + } injector; + + struct + { + _Atomic uint64_t jobs; + _Atomic uint64_t active; + _Atomic uint64_t max_active; + } read; + + struct + { + _Atomic uint64_t count[METADATA_OP_LATENCY_NKINDS]; + _Atomic uint64_t sum_ns[METADATA_OP_LATENCY_NKINDS]; + _Atomic uint64_t max_ns[METADATA_OP_LATENCY_NKINDS]; + _Atomic uint64_t buckets[METADATA_OP_LATENCY_NKINDS] + [METADATA_OP_LATENCY_NBUCKETS]; + } op_latency; +}; + +struct metadata_job +{ + struct metadata_job* next; + enum request_kind kind; + char* key; + uint64_t offset; + size_t requested_len; + metadata_store_read_cb read_cb; + metadata_store_stat_cb stat_cb; + void* user; +}; +struct metadata_store_async +{ + pthread_t* workers; + int started; + pthread_mutex_t lock; + pthread_cond_t wake; + struct metadata_job *pending_head, *pending_tail; + int stopping; + struct damacy_latency_model latency; + int latency_enabled; + uint64_t rng_state; + pthread_mutex_t rng_lock; + struct metadata_metrics metrics; +}; + +static void +atomic_max_u64(_Atomic uint64_t* dst, uint64_t val) +{ + uint64_t cur = atomic_load_explicit(dst, memory_order_relaxed); + while (cur < val && + !atomic_compare_exchange_weak_explicit( + dst, &cur, val, memory_order_relaxed, memory_order_relaxed)) { + } +} + +static uint64_t +monotonic_ns(void) +{ + struct timespec now; + clock_gettime(CLOCK_MONOTONIC, &now); + return (uint64_t)now.tv_sec * 1000000000ull + (uint64_t)now.tv_nsec; +} + +static unsigned +latency_bucket(uint64_t ns) +{ + if (ns == 0) + return 0; + unsigned idx = 63u - (unsigned)__builtin_clzll(ns); + if (idx >= METADATA_OP_LATENCY_NBUCKETS) + idx = METADATA_OP_LATENCY_NBUCKETS - 1; + return idx; +} + +static void +record_op_latency(struct metadata_store_async* s, + enum op_kind kind, + uint64_t submit_ts) +{ + if (!submit_ts) + return; + uint64_t now = monotonic_ns(); + uint64_t elapsed = now > submit_ts ? now - submit_ts : 0; + atomic_fetch_add_explicit( + &s->metrics.op_latency.count[kind], 1, memory_order_relaxed); + atomic_fetch_add_explicit( + &s->metrics.op_latency.sum_ns[kind], elapsed, memory_order_relaxed); + atomic_max_u64(&s->metrics.op_latency.max_ns[kind], elapsed); + atomic_fetch_add_explicit( + &s->metrics.op_latency.buckets[kind][latency_bucket(elapsed)], + 1, + memory_order_relaxed); +} + +static int +latency_enabled(const struct damacy_latency_model* l) +{ + return l && (l->baseline_ns || l->lognormal_mu_ln_ns != 0.0 || + l->lognormal_sigma_ln_ns != 0.0); +} + +static uint32_t +pcg32(uint64_t* state) +{ + uint64_t oldstate = *state; + *state = oldstate * 6364136223846793005ULL + 1442695040888963407ULL; + uint32_t xorshifted = (uint32_t)(((oldstate >> 18u) ^ oldstate) >> 27u); + uint32_t rot = (uint32_t)(oldstate >> 59u); + return (xorshifted >> rot) | (xorshifted << ((-rot) & 31u)); +} + +static double +uniform01(uint64_t* state) +{ + uint32_t u = pcg32(state); + return ((double)u + 1.0) / 4294967297.0; +} + +static double +normal01(uint64_t* state) +{ + double u1 = uniform01(state); + double u2 = uniform01(state); + return sqrt(-2.0 * log(u1)) * cos(6.2831853071795864769 * u2); +} + +static uint64_t +sample_delay_ns(struct metadata_store_async* s) +{ + const struct damacy_latency_model* l = &s->latency; + uint64_t ns = l->baseline_ns; + if (l->lognormal_mu_ln_ns == 0.0 && l->lognormal_sigma_ln_ns == 0.0) + return ns; + + pthread_mutex_lock(&s->rng_lock); + double z = normal01(&s->rng_state); + pthread_mutex_unlock(&s->rng_lock); + + double tail = exp(l->lognormal_mu_ln_ns + l->lognormal_sigma_ln_ns * z); + uint64_t tail_ns = tail > 0.0 ? (uint64_t)tail : 0; + if (l->cap_ns && tail_ns > l->cap_ns) + tail_ns = l->cap_ns; + if (UINT64_MAX - ns < tail_ns) + return UINT64_MAX; + return ns + tail_ns; +} + +static void +sleep_for_sample(struct metadata_store_async* s, enum latency_op_kind kind) +{ + if (!s->latency_enabled) + return; + uint64_t ns = sample_delay_ns(s); + if (ns > INT64_MAX) + ns = INT64_MAX; + atomic_fetch_add_explicit(&s->metrics.injector.ops, 1, memory_order_relaxed); + atomic_fetch_add_explicit( + &s->metrics.injector.total_sleep_ns, ns, memory_order_relaxed); + atomic_max_u64(&s->metrics.injector.max_sleep_ns, ns); + switch (kind) { + case LATENCY_OP_STAT: + atomic_fetch_add_explicit( + &s->metrics.injector.stat_ops, 1, memory_order_relaxed); + break; + case LATENCY_OP_SUBMIT: + atomic_fetch_add_explicit( + &s->metrics.injector.submit_ops, 1, memory_order_relaxed); + break; + } + uint64_t active = atomic_fetch_add_explicit( + &s->metrics.injector.active, 1, memory_order_relaxed) + + 1; + atomic_max_u64(&s->metrics.injector.max_active, active); + if (ns) + platform_sleep_ns((int64_t)ns); + atomic_fetch_sub_explicit( + &s->metrics.injector.active, 1, memory_order_relaxed); +} + +static void +read_active_begin(struct metadata_store_async* s) +{ + atomic_fetch_add_explicit(&s->metrics.read.jobs, 1, memory_order_relaxed); + uint64_t active = atomic_fetch_add_explicit( + &s->metrics.read.active, 1, memory_order_relaxed) + + 1; + atomic_max_u64(&s->metrics.read.max_active, active); +} + +static void +read_active_end(struct metadata_store_async* s) +{ + atomic_fetch_sub_explicit(&s->metrics.read.active, 1, memory_order_relaxed); +} + +static int +status_not_found_errno(int err) +{ + return err == ENOENT || err == ENOTDIR; +} + +static enum damacy_status +damacy_status_from_errno(int err) +{ + return status_not_found_errno(err) ? DAMACY_NOTFOUND : DAMACY_IO; +} + +static enum store_stat_result +stat_status_from_errno(int err) +{ + return status_not_found_errno(err) ? STORE_STAT_NOT_FOUND : STORE_STAT_ERROR; +} + +static void +job_free(struct metadata_job* job) +{ + free(job->key); + free(job); +} + +static struct metadata_job* +job_new(const char* key, enum request_kind kind) +{ + struct metadata_job* job = calloc(1, sizeof(*job)); + if (!job) + return NULL; + job->key = strdup(key); + if (!job->key) { + free(job); + return NULL; + } + job->kind = kind; + return job; +} + +static void +process_job(struct metadata_store_async* s, struct metadata_job* job) +{ + sleep_for_sample(s, + job->kind == REQ_STAT ? LATENCY_OP_STAT : LATENCY_OP_SUBMIT); + if (job->kind == REQ_STAT) { + struct stat st; + uint64_t started = monotonic_ns(); + int rc = stat(job->key, &st), error = errno; + record_op_latency(s, OP_STATX, started); + job->stat_cb(job->user, + rc ? stat_status_from_errno(error) : STORE_STAT_OK, + rc ? 0 : (uint64_t)st.st_size); + return; + } + enum damacy_status status = DAMACY_OK; + void* data = NULL; + size_t len = job->requested_len; + uint64_t started = monotonic_ns(); + int fd = open(job->key, O_RDONLY | O_CLOEXEC), error = errno; + record_op_latency(s, OP_OPEN, started); + if (fd < 0) { + status = damacy_status_from_errno(error); + goto Complete; + } + if (job->kind == REQ_READ_FILE) { + struct stat st; + started = monotonic_ns(); + int rc = fstat(fd, &st); + record_op_latency(s, OP_STATX, started); + if (rc || st.st_size < 0 || (uint64_t)st.st_size > SIZE_MAX) { + status = DAMACY_IO; + goto Close; + } + len = (size_t)st.st_size; + } + if (!len) + goto Close; + if (job->offset > INT64_MAX || len > (uint64_t)INT64_MAX - job->offset) { + status = DAMACY_IO; + goto Close; + } + data = malloc(len); + if (!data) { + status = DAMACY_OOM; + goto Close; + } + read_active_begin(s); + started = monotonic_ns(); + size_t done = 0; + while (done < len) { + ssize_t n = + pread(fd, (char*)data + done, len - done, (off_t)(job->offset + done)); + if (n < 0 && errno == EINTR) + continue; + if (n <= 0) { + status = DAMACY_IO; + break; + } + done += (size_t)n; + } + record_op_latency(s, OP_READ, started); + read_active_end(s); +Close: + started = monotonic_ns(); + if (close(fd) && status == DAMACY_OK) + log_warn("metadata_store_async: close failed for %s: %s", + job->key, + strerror(errno)); + record_op_latency(s, OP_CLOSE, started); +Complete: + if (status != DAMACY_OK) { + free(data); + data = NULL; + len = 0; + } + job->read_cb(job->user, status, data, len); +} + +static void* +worker_main(void* arg) +{ + struct metadata_store_async* s = arg; + for (;;) { + pthread_mutex_lock(&s->lock); + while (!s->pending_head && !s->stopping) + pthread_cond_wait(&s->wake, &s->lock); + struct metadata_job* job = s->pending_head; + if (!job) { + pthread_mutex_unlock(&s->lock); + break; + } + s->pending_head = job->next; + if (!s->pending_head) + s->pending_tail = NULL; + pthread_mutex_unlock(&s->lock); + process_job(s, job); + job_free(job); + } + return NULL; +} + +void +metadata_store_async_destroy(struct metadata_store_async* s) +{ + if (!s) + return; + pthread_mutex_lock(&s->lock); + s->stopping = 1; + pthread_cond_broadcast(&s->wake); + pthread_mutex_unlock(&s->lock); + for (int i = 0; i < s->started; ++i) + pthread_join(s->workers[i], NULL); + pthread_mutex_destroy(&s->rng_lock); + pthread_cond_destroy(&s->wake); + pthread_mutex_destroy(&s->lock); + free(s->workers); + free(s); +} + +struct metadata_store_async* +metadata_store_async_create(int concurrency, + const struct numa_resolved* affinity, + const struct damacy_latency_model* latency) +{ + (void)affinity; + if (concurrency < 1 || + (unsigned)concurrency > DAMACY_MAX_METADATA_IO_CONCURRENCY || + (latency && (!isfinite(latency->lognormal_mu_ln_ns) || + !isfinite(latency->lognormal_sigma_ln_ns) || + latency->lognormal_sigma_ln_ns < 0))) + return NULL; + struct metadata_store_async* s = calloc(1, sizeof(*s)); + if (!s) + return NULL; + s->workers = calloc((size_t)concurrency, sizeof(*s->workers)); + if (!s->workers) { + free(s); + return NULL; + } + if (pthread_mutex_init(&s->lock, NULL)) + goto Fail; + if (pthread_cond_init(&s->wake, NULL)) { + pthread_mutex_destroy(&s->lock); + goto Fail; + } + if (pthread_mutex_init(&s->rng_lock, NULL)) { + pthread_cond_destroy(&s->wake); + pthread_mutex_destroy(&s->lock); + goto Fail; + } + if (latency) + s->latency = *latency; + s->latency_enabled = latency_enabled(latency); + s->rng_state = s->latency.seed ? s->latency.seed : 0xc0ffee1234ULL; + for (int i = 0; i < concurrency; ++i) { + if (pthread_create(&s->workers[i], NULL, worker_main, s)) { + metadata_store_async_destroy(s); + return NULL; + } + ++s->started; + } + return s; +Fail: + free(s->workers); + free(s); + return NULL; +} + +static int +enqueue_job(struct metadata_store_async* s, struct metadata_job* job) +{ + pthread_mutex_lock(&s->lock); + if (s->stopping) { + pthread_mutex_unlock(&s->lock); + return 1; + } + if (s->pending_tail) + s->pending_tail->next = job; + else + s->pending_head = job; + s->pending_tail = job; + pthread_cond_signal(&s->wake); + pthread_mutex_unlock(&s->lock); + return 0; +} + +static int +post_read(struct metadata_store_async* s, + const char* key, + uint64_t offset, + size_t len, + enum request_kind kind, + metadata_store_read_cb cb, + void* user) +{ + if (!s || !key || !cb) + return 1; + struct metadata_job* job = job_new(key, kind); + if (!job) + return 1; + job->offset = offset; + job->requested_len = len; + job->read_cb = cb; + job->user = user; + if (enqueue_job(s, job)) { + job_free(job); + return 1; + } + return 0; +} + +int +metadata_store_async_read_file(struct metadata_store_async* s, + const char* key, + metadata_store_read_cb cb, + void* user) +{ + return post_read(s, key, 0, 0, REQ_READ_FILE, cb, user); +} + +int +metadata_store_async_read(struct metadata_store_async* s, + const char* key, + uint64_t offset, + size_t len, + metadata_store_read_cb cb, + void* user) +{ + return post_read(s, key, offset, len, REQ_READ, cb, user); +} + +int +metadata_store_async_stat(struct metadata_store_async* s, + const char* key, + metadata_store_stat_cb cb, + void* user) +{ + if (!s || !key || !cb) + return 1; + struct metadata_job* job = job_new(key, REQ_STAT); + if (!job) + return 1; + job->stat_cb = cb; + job->user = user; + if (enqueue_job(s, job)) { + job_free(job); + return 1; + } + return 0; +} + +void +metadata_store_async_latency_stats_get( + struct metadata_store_async* s, + struct metadata_store_async_latency_stats* out) +{ + if (!out) + return; + *out = (struct metadata_store_async_latency_stats){ 0 }; + if (!s) + return; + *out = (struct metadata_store_async_latency_stats){ + .ops = atomic_load_explicit(&s->metrics.injector.ops, memory_order_relaxed), + .stat_ops = + atomic_load_explicit(&s->metrics.injector.stat_ops, memory_order_relaxed), + .submit_ops = atomic_load_explicit(&s->metrics.injector.submit_ops, + memory_order_relaxed), + .active = + atomic_load_explicit(&s->metrics.injector.active, memory_order_relaxed), + .max_active = atomic_load_explicit(&s->metrics.injector.max_active, + memory_order_relaxed), + .total_sleep_ns = atomic_load_explicit(&s->metrics.injector.total_sleep_ns, + memory_order_relaxed), + .max_sleep_ns = atomic_load_explicit(&s->metrics.injector.max_sleep_ns, + memory_order_relaxed), + }; +} + +void +metadata_store_async_latency_stats_reset(struct metadata_store_async* s) +{ + if (!s) + return; + uint64_t active = + atomic_load_explicit(&s->metrics.injector.active, memory_order_relaxed); + atomic_store_explicit(&s->metrics.injector.ops, 0, memory_order_relaxed); + atomic_store_explicit(&s->metrics.injector.stat_ops, 0, memory_order_relaxed); + atomic_store_explicit( + &s->metrics.injector.submit_ops, 0, memory_order_relaxed); + atomic_store_explicit( + &s->metrics.injector.max_active, active, memory_order_relaxed); + atomic_store_explicit( + &s->metrics.injector.total_sleep_ns, 0, memory_order_relaxed); + atomic_store_explicit( + &s->metrics.injector.max_sleep_ns, 0, memory_order_relaxed); +} + +void +metadata_store_async_backend_stats_get( + struct metadata_store_async* s, + struct metadata_store_async_backend_stats* out) +{ + if (!out) + return; + *out = (struct metadata_store_async_backend_stats){ 0 }; + if (!s) + return; + *out = (struct metadata_store_async_backend_stats){ + .read_jobs = + atomic_load_explicit(&s->metrics.read.jobs, memory_order_relaxed), + .read_active = + atomic_load_explicit(&s->metrics.read.active, memory_order_relaxed), + .read_max_active = + atomic_load_explicit(&s->metrics.read.max_active, memory_order_relaxed), + }; +} + +void +metadata_store_async_backend_stats_reset(struct metadata_store_async* s) +{ + if (!s) + return; + uint64_t active = + atomic_load_explicit(&s->metrics.read.active, memory_order_relaxed); + atomic_store_explicit(&s->metrics.read.jobs, 0, memory_order_relaxed); + atomic_store_explicit( + &s->metrics.read.max_active, active, memory_order_relaxed); +} + +void +metadata_store_async_op_latency_stats_get( + struct metadata_store_async* s, + struct metadata_store_async_op_latency_stats* out) +{ + if (!out) + return; + *out = (struct metadata_store_async_op_latency_stats){ 0 }; + if (!s) + return; + for (unsigned k = 0; k < METADATA_OP_LATENCY_NKINDS; ++k) { + out->kinds[k].count = atomic_load_explicit(&s->metrics.op_latency.count[k], + memory_order_relaxed); + out->kinds[k].sum_ns = atomic_load_explicit( + &s->metrics.op_latency.sum_ns[k], memory_order_relaxed); + out->kinds[k].max_ns = atomic_load_explicit( + &s->metrics.op_latency.max_ns[k], memory_order_relaxed); + for (unsigned b = 0; b < METADATA_OP_LATENCY_NBUCKETS; ++b) + out->kinds[k].buckets[b] = atomic_load_explicit( + &s->metrics.op_latency.buckets[k][b], memory_order_relaxed); + } +} + +void +metadata_store_async_op_latency_stats_reset(struct metadata_store_async* s) +{ + if (!s) + return; + for (unsigned k = 0; k < METADATA_OP_LATENCY_NKINDS; ++k) { + atomic_store_explicit( + &s->metrics.op_latency.count[k], 0, memory_order_relaxed); + atomic_store_explicit( + &s->metrics.op_latency.sum_ns[k], 0, memory_order_relaxed); + atomic_store_explicit( + &s->metrics.op_latency.max_ns[k], 0, memory_order_relaxed); + for (unsigned b = 0; b < METADATA_OP_LATENCY_NBUCKETS; ++b) + atomic_store_explicit( + &s->metrics.op_latency.buckets[k][b], 0, memory_order_relaxed); + } +} diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index f4b0c6fe..3d9c758c 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -30,6 +30,10 @@ add_damacy_test(test_io_queue io_queue log platform) add_damacy_test(test_store_fs store lru hash log) add_damacy_test(test_store_latency store test_fixture) add_damacy_test(test_metadata_store_async metadata_store_async test_fixture) +set_tests_properties(test_metadata_store_async PROPERTIES TIMEOUT 30) +if(APPLE) + add_damacy_test(test_platform_mach platform log) +endif() target_link_libraries(test_store_fs PRIVATE Threads::Threads) target_link_libraries(test_metadata_store_async PRIVATE Threads::Threads) add_damacy_test( diff --git a/tests/test_array_meta.c b/tests/test_array_meta.c index d287766d..70c3d94d 100644 --- a/tests/test_array_meta.c +++ b/tests/test_array_meta.c @@ -12,6 +12,7 @@ #include #include #include +#include static const char* MINIMAL_ZARR_JSON = "{" diff --git a/tests/test_metadata_store_async.c b/tests/test_metadata_store_async.c index ed1f32f6..bf95f640 100644 --- a/tests/test_metadata_store_async.c +++ b/tests/test_metadata_store_async.c @@ -309,11 +309,151 @@ test_read_file_open_error_completes(void) return 0; } +static int +test_empty_and_short_reads(void) +{ + char root[128], path[256]; + EXPECT(make_root(root, sizeof root) == 0); + snprintf(path, sizeof path, "%s/payload", root); + EXPECT(fixture_write_file(path, "") == 0); + struct metadata_store_async* s = metadata_store_async_create(2, NULL, NULL); + EXPECT(s); + + struct read_wait empty, zero, short_read; + read_wait_init(&empty); + read_wait_init(&zero); + read_wait_init(&short_read); + EXPECT(metadata_store_async_read_file(s, path, read_cb, &empty) == 0); + EXPECT(metadata_store_async_read(s, path, 17, 0, read_cb, &zero) == 0); + read_wait_block(&empty); + read_wait_block(&zero); + EXPECT(empty.status == DAMACY_OK && empty.len == 0 && !empty.data); + EXPECT(zero.status == DAMACY_OK && zero.len == 0 && !zero.data); + struct metadata_store_async_backend_stats stats; + metadata_store_async_backend_stats_get(s, &stats); + EXPECT(stats.read_jobs == 0 && stats.read_active == 0); + + EXPECT(fixture_write_file(path, "abc") == 0); + EXPECT(metadata_store_async_read(s, path, 1, 3, read_cb, &short_read) == 0); + read_wait_block(&short_read); + EXPECT(short_read.status == DAMACY_IO); + EXPECT(short_read.len == 0 && !short_read.data); + metadata_store_async_destroy(s); + read_wait_done(&empty); + read_wait_done(&zero); + read_wait_done(&short_read); + fixture_rm_tree(root); + return 0; +} + +static int +test_missing_stat_and_not_directory(void) +{ + char root[128], path[256], missing[300]; + EXPECT(make_root(root, sizeof root) == 0); + snprintf(path, sizeof path, "%s/payload", root); + EXPECT(fixture_write_file(path, "abc") == 0); + struct metadata_store_async* s = metadata_store_async_create(2, NULL, NULL); + EXPECT(s); + for (int not_directory = 0; not_directory < 2; ++not_directory) { + snprintf( + missing, sizeof missing, "%s/missing", not_directory ? path : root); + struct read_wait r; + struct stat_wait st; + read_wait_init(&r); + stat_wait_init(&st); + EXPECT(metadata_store_async_read_file(s, missing, read_cb, &r) == 0); + EXPECT(metadata_store_async_stat(s, missing, stat_cb, &st) == 0); + read_wait_block(&r); + stat_wait_block(&st); + EXPECT(r.status == DAMACY_NOTFOUND && !r.data && r.len == 0); + EXPECT(st.status == STORE_STAT_NOT_FOUND && st.size == 0); + read_wait_done(&r); + stat_wait_done(&st); + } + metadata_store_async_destroy(s); + fixture_rm_tree(root); + return 0; +} + +struct chain_read +{ + struct metadata_store_async* store; + const char* path; + struct read_wait first, next; + int submit_failed; +}; + +static void +chain_read_cb(void* user, enum damacy_status status, void* data, size_t len) +{ + struct chain_read* chain = user; + chain->submit_failed = metadata_store_async_read( + chain->store, chain->path, 1, 2, read_cb, &chain->next); + read_cb(&chain->first, status, data, len); +} + +static int +test_callback_can_submit(void) +{ + char root[128], path[256]; + EXPECT(make_root(root, sizeof root) == 0); + snprintf(path, sizeof path, "%s/payload", root); + EXPECT(fixture_write_file(path, "abc") == 0); + struct metadata_store_async* s = metadata_store_async_create(1, NULL, NULL); + EXPECT(s); + struct chain_read chain = { .store = s, .path = path }; + read_wait_init(&chain.first); + read_wait_init(&chain.next); + EXPECT(metadata_store_async_read_file(s, path, chain_read_cb, &chain) == 0); + EXPECT(read_wait_block_ms(&chain.first, 2000)); + EXPECT(!chain.submit_failed); + EXPECT(read_wait_block_ms(&chain.next, 2000)); + EXPECT(chain.first.status == DAMACY_OK && chain.first.len == 3); + EXPECT(chain.next.status == DAMACY_OK && chain.next.len == 2); + EXPECT(!memcmp(chain.next.data, "bc", 2)); + metadata_store_async_destroy(s); + read_wait_done(&chain.first); + read_wait_done(&chain.next); + fixture_rm_tree(root); + return 0; +} + +static int +test_destroy_drains_requests(void) +{ + char root[128], path[256]; + EXPECT(make_root(root, sizeof root) == 0); + snprintf(path, sizeof path, "%s/payload", root); + EXPECT(fixture_write_file(path, "abc") == 0); + struct metadata_store_async* s = metadata_store_async_create( + 2, NULL, &(struct damacy_latency_model){ .baseline_ns = 1000000 }); + EXPECT(s); + struct read_wait pending[32]; + for (size_t i = 0; i < sizeof pending / sizeof *pending; ++i) { + read_wait_init(&pending[i]); + EXPECT(metadata_store_async_read_file(s, path, read_cb, &pending[i]) == 0); + } + metadata_store_async_destroy(s); + for (size_t i = 0; i < sizeof pending / sizeof *pending; ++i) { + EXPECT(pending[i].done); + EXPECT(pending[i].status == DAMACY_OK && pending[i].len == 3); + EXPECT(!memcmp(pending[i].data, "abc", 3)); + read_wait_done(&pending[i]); + } + fixture_rm_tree(root); + return 0; +} + int main(void) { RUN(test_read_stat_and_stats); RUN(test_op_latency_measured); RUN(test_read_file_open_error_completes); + RUN(test_empty_and_short_reads); + RUN(test_missing_stat_and_not_directory); + RUN(test_callback_can_submit); + RUN(test_destroy_drains_requests); return 0; } diff --git a/tests/test_platform_mach.c b/tests/test_platform_mach.c new file mode 100644 index 00000000..4b804ed9 --- /dev/null +++ b/tests/test_platform_mach.c @@ -0,0 +1,47 @@ +#include "expect.h" +#include "platform/numa.h" +#include "platform/platform.h" + +#include +#include +#include + +static int +test_memory_and_threads(void) +{ + EXPECT(platform_page_size() == (size_t)sysconf(_SC_PAGESIZE)); + EXPECT(platform_page_alignment() == platform_page_size()); + EXPECT(platform_available_memory() > 0); + EXPECT(platform_default_thread_count() > 0); + void* ptr = platform_aligned_alloc(64, 128); + EXPECT(ptr && (uintptr_t)ptr % 64 == 0); + platform_aligned_free(ptr); + return 0; +} + +static int +test_unavailable_affinity(void) +{ + EXPECT(!platform_numa_available()); + EXPECT(platform_numa_max_node() == -1); + struct platform_cpu_mask mask; + memset(&mask, 0xff, sizeof mask); + EXPECT(platform_numa_node_cpu_mask(0, &mask) != 0); + EXPECT(platform_cpu_mask_is_empty(&mask)); + memset(&mask, 0xff, sizeof mask); + EXPECT(platform_thread_affinity_get(&mask) != 0); + EXPECT(platform_cpu_mask_is_empty(&mask)); + EXPECT(platform_thread_affinity_set(&mask) != 0); + int first = 0, last = 0, count = 1; + EXPECT(platform_cpu_mask_describe(&mask, &first, &last, &count) != 0); + EXPECT(first == -1 && last == -1 && count == 0); + return 0; +} + +int +main(void) +{ + RUN(test_memory_and_threads); + RUN(test_unavailable_affinity); + return 0; +} diff --git a/tests/test_store_latency.c b/tests/test_store_latency.c index 6386bd65..3e080cc1 100644 --- a/tests/test_store_latency.c +++ b/tests/test_store_latency.c @@ -10,6 +10,7 @@ #include #include #include +#include static uint64_t now_ns(void) From 42548fc00b49fccb91c0602a62201871174b9cc4 Mon Sep 17 00:00:00 2001 From: Nathan Clack Date: Wed, 23 Sep 2026 17:16:51 -0700 Subject: [PATCH 2/9] Run I/O queue ordering test on small CI runners --- tests/test_io_queue.c | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/tests/test_io_queue.c b/tests/test_io_queue.c index 1df18eba..0cb748c9 100644 --- a/tests/test_io_queue.c +++ b/tests/test_io_queue.c @@ -140,10 +140,11 @@ test_out_of_order_completion(void) // Regression: retired_seq must track the lowest in-flight seq even // when higher seqs finish first. Posts W gate jobs, releases all but // gate[0], waits for a barrier post — its retirement proves the - // higher seqs completed without advancing past gate[0]. + // higher seqs completed without advancing past gate[0]. Two workers are + // enough to exercise this invariant and fit small hosted CI runners. enum { - W = 4 + W = 2 }; struct io_queue* q = io_queue_create(W, NULL); EXPECT(q); From 859045c76293a2820ca2d0bf3564ecc7cfc3d61f Mon Sep 17 00:00:00 2001 From: Nathan Clack Date: Thu, 24 Sep 2026 21:16:52 +0000 Subject: [PATCH 3/9] platform: build posix sources on macOS --- src/platform/platform.h | 1 - src/platform/platform.mach.c | 295 ---------------------------------- src/platform/platform.posix.c | 10 -- tests/test_platform_mach.c | 1 - 4 files changed, 307 deletions(-) delete mode 100644 src/platform/platform.mach.c diff --git a/src/platform/platform.h b/src/platform/platform.h index 0380d56c..58f1af79 100644 --- a/src/platform/platform.h +++ b/src/platform/platform.h @@ -14,7 +14,6 @@ extern "C" size_t platform_page_alignment(void); void* platform_aligned_alloc(size_t alignment, size_t size); void platform_aligned_free(void* ptr); - size_t platform_available_memory(void); struct platform_clock { diff --git a/src/platform/platform.mach.c b/src/platform/platform.mach.c deleted file mode 100644 index 15a08aae..00000000 --- a/src/platform/platform.mach.c +++ /dev/null @@ -1,295 +0,0 @@ -#include "platform/platform.h" - -#include -#include -#include -#include -#include -#include -#include -#include -#include - -size_t -platform_page_size(void) -{ - long size = sysconf(_SC_PAGESIZE); - return size > 0 ? (size_t)size : 0; -} - -size_t -platform_page_alignment(void) -{ - size_t ps = platform_page_size(); - return ps > 0 ? ps : 4096; -} - -size_t -platform_available_memory(void) -{ - vm_statistics64_data_t vm; - mach_msg_type_number_t count = HOST_VM_INFO64_COUNT; - mach_port_t host = mach_host_self(); - kern_return_t rc = - host_statistics64(host, HOST_VM_INFO64, (host_info64_t)&vm, &count); - mach_port_deallocate(mach_task_self(), host); - if (rc != KERN_SUCCESS) - return 0; - return ((size_t)vm.free_count + vm.inactive_count) * platform_page_size(); -} - -void* -platform_aligned_alloc(size_t alignment, size_t size) -{ - return aligned_alloc(alignment, size); -} - -void -platform_aligned_free(void* ptr) -{ - free(ptr); -} - -void -platform_sleep_ns(int64_t ns) -{ - struct timespec ts = { - .tv_sec = ns / 1000000000LL, - .tv_nsec = ns % 1000000000LL, - }; - nanosleep(&ts, NULL); -} - -static int64_t -monotonic_ns(void) -{ - struct timespec now; - clock_gettime(CLOCK_MONOTONIC, &now); - return (int64_t)now.tv_sec * 1000000000LL + now.tv_nsec; -} - -float -platform_toc(struct platform_clock* clock) -{ - int64_t now = monotonic_ns(); - float elapsed = (float)((now - clock->last_ns) / 1e9); - clock->last_ns = now; - return elapsed; -} - -void -platform_localtime(const time_t* t, struct tm* out) -{ - localtime_r(t, out); -} - -void -platform_call_once(platform_once* flag, void (*fn)(void)) -{ - pthread_once(flag, fn); -} - -struct platform_thread -{ - pthread_t handle; - void (*fn)(void*); - void* arg; -}; - -static void* -thread_trampoline(void* p) -{ - struct platform_thread* t = (struct platform_thread*)p; - t->fn(t->arg); - return NULL; -} - -struct platform_thread* -platform_thread_start(void (*fn)(void*), void* arg) -{ - struct platform_thread* t = - (struct platform_thread*)calloc(1, sizeof(struct platform_thread)); - if (!t) - return NULL; - t->fn = fn; - t->arg = arg; - if (pthread_create(&t->handle, NULL, thread_trampoline, t) != 0) { - free(t); - return NULL; - } - return t; -} - -int -platform_thread_join(struct platform_thread* t) -{ - if (!t) - return -1; - int rc = pthread_join(t->handle, NULL); - free(t); - return rc; -} - -struct platform_mutex -{ - pthread_mutex_t m; -}; - -struct platform_mutex* -platform_mutex_new(void) -{ - struct platform_mutex* m = - (struct platform_mutex*)calloc(1, sizeof(struct platform_mutex)); - if (!m) - return NULL; - if (pthread_mutex_init(&m->m, NULL) != 0) { - free(m); - return NULL; - } - return m; -} - -void -platform_mutex_free(struct platform_mutex* m) -{ - if (!m) - return; - pthread_mutex_destroy(&m->m); - free(m); -} - -void -platform_mutex_lock(struct platform_mutex* m) -{ - pthread_mutex_lock(&m->m); -} - -void -platform_mutex_unlock(struct platform_mutex* m) -{ - pthread_mutex_unlock(&m->m); -} - -struct platform_cond -{ - pthread_cond_t c; -}; - -struct platform_cond* -platform_cond_new(void) -{ - struct platform_cond* c = - (struct platform_cond*)calloc(1, sizeof(struct platform_cond)); - if (!c) - return NULL; - if (pthread_cond_init(&c->c, NULL) != 0) { - free(c); - return NULL; - } - return c; -} - -void -platform_cond_free(struct platform_cond* c) -{ - if (!c) - return; - pthread_cond_destroy(&c->c); - free(c); -} - -void -platform_cond_wait(struct platform_cond* c, struct platform_mutex* m) -{ - pthread_cond_wait(&c->c, &m->m); -} - -int -platform_cond_timedwait_ms(struct platform_cond* c, - struct platform_mutex* m, - int timeout_ms) -{ - struct timespec ts; - clock_gettime(CLOCK_REALTIME, &ts); - ts.tv_sec += timeout_ms / 1000; - ts.tv_nsec += (long)(timeout_ms % 1000) * 1000000L; - if (ts.tv_nsec >= 1000000000L) { - ts.tv_sec += 1; - ts.tv_nsec -= 1000000000L; - } - return pthread_cond_timedwait(&c->c, &m->m, &ts) == ETIMEDOUT ? 1 : 0; -} - -void -platform_cond_broadcast(struct platform_cond* c) -{ - pthread_cond_broadcast(&c->c); -} - -void -platform_cpu_pause(void) -{ -#if defined(__i386__) || defined(__x86_64__) - __asm__ __volatile__("pause" ::: "memory"); -#elif defined(__aarch64__) || defined(__arm__) - __asm__ __volatile__("yield" ::: "memory"); -#else - sched_yield(); -#endif -} - -void -platform_yield(void) -{ - sched_yield(); -} - -int -platform_default_thread_count(void) -{ - long n = sysconf(_SC_NPROCESSORS_ONLN); - return n > 0 ? (int)n : 1; -} - -void* -platform_dlopen(const char* path) -{ - return dlopen(path, RTLD_LAZY | RTLD_LOCAL); -} - -void* -platform_dlsym(void* handle, const char* name) -{ - return dlsym(handle, name); -} - -void -platform_dlclose(void* handle) -{ - if (!handle) - return; - dlclose(handle); -} - -const char* -platform_dlerror(void) -{ - return dlerror(); -} - -const char* -platform_getenv(const char* name) -{ - return getenv(name); -} - -uint64_t -platform_max_open_files(void) -{ - struct rlimit r; - if (getrlimit(RLIMIT_NOFILE, &r) != 0) - return 0; - if (r.rlim_cur == RLIM_INFINITY) - return 0; - return (uint64_t)r.rlim_cur; -} diff --git a/src/platform/platform.posix.c b/src/platform/platform.posix.c index c642e809..e7fd60b4 100644 --- a/src/platform/platform.posix.c +++ b/src/platform/platform.posix.c @@ -22,16 +22,6 @@ platform_page_alignment(void) return ps > 0 ? ps : 4096; } -size_t -platform_available_memory(void) -{ - long pages = sysconf(_SC_AVPHYS_PAGES); - long page_sz = sysconf(_SC_PAGESIZE); - if (pages > 0 && page_sz > 0) - return (size_t)pages * (size_t)page_sz; - return 0; -} - void* platform_aligned_alloc(size_t alignment, size_t size) { diff --git a/tests/test_platform_mach.c b/tests/test_platform_mach.c index 4b804ed9..fd5a8610 100644 --- a/tests/test_platform_mach.c +++ b/tests/test_platform_mach.c @@ -11,7 +11,6 @@ test_memory_and_threads(void) { EXPECT(platform_page_size() == (size_t)sysconf(_SC_PAGESIZE)); EXPECT(platform_page_alignment() == platform_page_size()); - EXPECT(platform_available_memory() > 0); EXPECT(platform_default_thread_count() > 0); void* ptr = platform_aligned_alloc(64, 128); EXPECT(ptr && (uintptr_t)ptr % 64 == 0); From 75adf8ffdb3e37c4d69b166e21550770363ae67e Mon Sep 17 00:00:00 2001 From: Nathan Clack Date: Thu, 24 Sep 2026 21:17:33 +0000 Subject: [PATCH 4/9] platform: name macOS NUMA stub .darwin.c --- cmake/Helpers.cmake | 5 +---- docs/pipeline.md | 4 ++-- src/CMakeLists.txt | 4 ++-- src/platform/{numa.mach.c => numa.darwin.c} | 0 tests/CMakeLists.txt | 2 +- tests/{test_platform_mach.c => test_platform_darwin.c} | 0 6 files changed, 6 insertions(+), 9 deletions(-) rename src/platform/{numa.mach.c => numa.darwin.c} (100%) rename tests/{test_platform_mach.c => test_platform_darwin.c} (100%) diff --git a/cmake/Helpers.cmake b/cmake/Helpers.cmake index 9648cfa4..d8b19731 100644 --- a/cmake/Helpers.cmake +++ b/cmake/Helpers.cmake @@ -31,8 +31,7 @@ function(add_cuda_lib TARGET) add_src_lib(${TARGET} ${ARGN}) endfunction() -# Appends basename.win32.c or basename.posix.c (.mach.c on macOS, when -# present, with a legacy .darwin.c fallback) +# Appends basename.win32.c or basename.posix.c (or .darwin.c if it exists) # to TARGET's sources. BASENAME defaults to TARGET. function(add_platform_sources TARGET) if(ARGC GREATER 1) @@ -42,8 +41,6 @@ function(add_platform_sources TARGET) endif() if(WIN32) target_sources(${TARGET} PRIVATE ${BASE}.win32.c) - elseif(APPLE AND EXISTS "${CMAKE_CURRENT_SOURCE_DIR}/${BASE}.mach.c") - target_sources(${TARGET} PRIVATE ${BASE}.mach.c) elseif(APPLE AND EXISTS "${CMAKE_CURRENT_SOURCE_DIR}/${BASE}.darwin.c") target_sources(${TARGET} PRIVATE ${BASE}.darwin.c) else() diff --git a/docs/pipeline.md b/docs/pipeline.md index 5bc81c28..8149ca2e 100644 --- a/docs/pipeline.md +++ b/docs/pipeline.md @@ -207,8 +207,8 @@ The extension has no CUDA or nvCOMP dependency in this configuration. `CudaExecutor` reports that CUDA support was not built. Linux retains its io_uring metadata requirements. macOS uses a POSIX metadata worker pool, with one worker per `metadata_io_concurrency`, and shared POSIX bulk reads. -CMake selects `platform.mach.c` and `numa.mach.c`; NUMA placement and CPU -affinity are unavailable. CUDA defaults off on macOS and cannot be enabled. +NUMA placement and CPU affinity are unavailable. CUDA defaults off on macOS +and cannot be enabled. On Linux CUDA defaults on, builds both executors, and requires the CUDA toolkit, nvCOMP, and a runtime NVIDIA driver. GDS requires a CUDA build. diff --git a/src/CMakeLists.txt b/src/CMakeLists.txt index 024d3e65..71f604cc 100644 --- a/src/CMakeLists.txt +++ b/src/CMakeLists.txt @@ -32,8 +32,8 @@ add_src_lib(lru SOURCES util/lru.c util/lru.h LINKS log) add_src_lib(pool SOURCES util/pool.c util/pool.h LINKS log platform) add_src_lib(path_intern SOURCES util/path_intern.c util/path_intern.h LINKS hash) -# Platform abstraction: CMake selects Mach implementations on macOS and -# POSIX implementations on Linux. Shared POSIX file I/O works on both. +# Platform abstraction: POSIX sources on Linux and macOS. macOS has no NUMA +# or CPU affinity controls, so it builds the numa.darwin.c stub instead. # platform.posix.c uses dlopen for platform_dl*; numa.posix.c uses libnuma via # platform_dl* so a single binary works on hosts with or without libnuma installed. add_src_lib(platform diff --git a/src/platform/numa.mach.c b/src/platform/numa.darwin.c similarity index 100% rename from src/platform/numa.mach.c rename to src/platform/numa.darwin.c diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index 3d9c758c..c2394f3f 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -32,7 +32,7 @@ add_damacy_test(test_store_latency store test_fixture) add_damacy_test(test_metadata_store_async metadata_store_async test_fixture) set_tests_properties(test_metadata_store_async PROPERTIES TIMEOUT 30) if(APPLE) - add_damacy_test(test_platform_mach platform log) + add_damacy_test(test_platform_darwin platform log) endif() target_link_libraries(test_store_fs PRIVATE Threads::Threads) target_link_libraries(test_metadata_store_async PRIVATE Threads::Threads) diff --git a/tests/test_platform_mach.c b/tests/test_platform_darwin.c similarity index 100% rename from tests/test_platform_mach.c rename to tests/test_platform_darwin.c From e1117349f80af4be5d44a0aa3657b8d54d8bfca0 Mon Sep 17 00:00:00 2001 From: Nathan Clack Date: Thu, 24 Sep 2026 21:17:40 +0000 Subject: [PATCH 5/9] store: metadata .mach.c -> .posix.c --- src/CMakeLists.txt | 2 +- ...metadata_store_async.mach.c => metadata_store_async.posix.c} | 0 2 files changed, 1 insertion(+), 1 deletion(-) rename src/store/{metadata_store_async.mach.c => metadata_store_async.posix.c} (100%) diff --git a/src/CMakeLists.txt b/src/CMakeLists.txt index 71f604cc..aefd6a9f 100644 --- a/src/CMakeLists.txt +++ b/src/CMakeLists.txt @@ -190,7 +190,7 @@ else() endif() if(APPLE) - set(METADATA_STORE_IMPL store/metadata_store_async.mach.c) + set(METADATA_STORE_IMPL store/metadata_store_async.posix.c) else() set(METADATA_STORE_IMPL store/metadata_store_async.c) endif() diff --git a/src/store/metadata_store_async.mach.c b/src/store/metadata_store_async.posix.c similarity index 100% rename from src/store/metadata_store_async.mach.c rename to src/store/metadata_store_async.posix.c From 1f6f004260e874056f6cf0b1b13f62b5d64c081b Mon Sep 17 00:00:00 2001 From: Nathan Clack Date: Thu, 24 Sep 2026 21:20:55 +0000 Subject: [PATCH 6/9] store: share metadata backend stats code --- src/CMakeLists.txt | 3 +- src/store/metadata_store_async.c | 379 ++--------------------- src/store/metadata_store_async.posix.c | 387 ++---------------------- src/store/metadata_store_async_common.c | 306 +++++++++++++++++++ src/store/metadata_store_async_common.h | 91 ++++++ 5 files changed, 448 insertions(+), 718 deletions(-) create mode 100644 src/store/metadata_store_async_common.c create mode 100644 src/store/metadata_store_async_common.h diff --git a/src/CMakeLists.txt b/src/CMakeLists.txt index aefd6a9f..27fe67fc 100644 --- a/src/CMakeLists.txt +++ b/src/CMakeLists.txt @@ -195,7 +195,8 @@ else() set(METADATA_STORE_IMPL store/metadata_store_async.c) endif() add_src_lib(metadata_store_async - SOURCES store/metadata_store_async.h ${METADATA_STORE_IMPL} + SOURCES store/metadata_store_async.h store/metadata_store_async_common.h + store/metadata_store_async_common.c ${METADATA_STORE_IMPL} LINKS store log platform ) if(NOT APPLE) diff --git a/src/store/metadata_store_async.c b/src/store/metadata_store_async.c index fb75fc45..80d34541 100644 --- a/src/store/metadata_store_async.c +++ b/src/store/metadata_store_async.c @@ -3,9 +3,9 @@ #endif #include "store/metadata_store_async.h" +#include "store/metadata_store_async_common.h" #include "log/log.h" -#include "platform/platform.h" #ifdef DAMACY_METADATA_ASYNC_NUMA #include "numa/numa.h" #endif @@ -16,14 +16,12 @@ #include #include #include -#include #include #include #include #include #include #include -#include #include #define URING_MAX_EPOLL_EVENTS 2 @@ -35,20 +33,6 @@ enum request_kind REQ_STAT, }; -enum op_kind -{ - OP_STATX, - OP_OPEN, - OP_READ, - OP_CLOSE, -}; - -enum latency_op_kind -{ - LATENCY_OP_STAT, - LATENCY_OP_SUBMIT, -}; - #define OP_TAG_MASK 0x7ull struct metadata_job; @@ -86,33 +70,6 @@ struct metadata_job _Static_assert(_Alignof(struct metadata_job) >= 8, "metadata_job pointers must leave low bits for op tags"); -struct metadata_metrics -{ - struct { - _Atomic uint64_t ops; - _Atomic uint64_t stat_ops; - _Atomic uint64_t submit_ops; - _Atomic uint64_t active; - _Atomic uint64_t max_active; - _Atomic uint64_t total_sleep_ns; - _Atomic uint64_t max_sleep_ns; - } injector; - - struct { - _Atomic uint64_t jobs; - _Atomic uint64_t active; - _Atomic uint64_t max_active; - } read; - - struct { - _Atomic uint64_t count[METADATA_OP_LATENCY_NKINDS]; - _Atomic uint64_t sum_ns[METADATA_OP_LATENCY_NKINDS]; - _Atomic uint64_t max_ns[METADATA_OP_LATENCY_NKINDS]; - _Atomic uint64_t buckets[METADATA_OP_LATENCY_NKINDS] - [METADATA_OP_LATENCY_NBUCKETS]; - } op_latency; -}; - struct metadata_store_async { struct io_uring ring; @@ -120,7 +77,7 @@ struct metadata_store_async int worker_started; int ring_initialized; int lock_initialized; - int rng_lock_initialized; + int common_initialized; int event_fd; int epoll_fd; uint32_t concurrency; @@ -131,185 +88,13 @@ struct metadata_store_async struct metadata_job* pending_tail; int stopping; - struct damacy_latency_model latency; - int latency_enabled; - uint64_t rng_state; - pthread_mutex_t rng_lock; - - struct metadata_metrics metrics; + struct metadata_store_common common; #ifdef DAMACY_METADATA_ASYNC_NUMA struct numa_resolved affinity; #endif }; -static void -atomic_max_u64(_Atomic uint64_t* dst, uint64_t val) -{ - uint64_t cur = atomic_load_explicit(dst, memory_order_relaxed); - while (cur < val && - !atomic_compare_exchange_weak_explicit( - dst, &cur, val, memory_order_relaxed, memory_order_relaxed)) { - } -} - -static uint64_t -monotonic_ns(void) -{ - struct timespec now; - clock_gettime(CLOCK_MONOTONIC, &now); - return (uint64_t)now.tv_sec * 1000000000ull + (uint64_t)now.tv_nsec; -} - -static unsigned -latency_bucket(uint64_t ns) -{ - if (ns == 0) - return 0; - unsigned idx = 63u - (unsigned)__builtin_clzll(ns); - if (idx >= METADATA_OP_LATENCY_NBUCKETS) - idx = METADATA_OP_LATENCY_NBUCKETS - 1; - return idx; -} - -static void -record_op_latency(struct metadata_store_async* s, - enum op_kind kind, - uint64_t submit_ts) -{ - if (!submit_ts) - return; - uint64_t now = monotonic_ns(); - uint64_t elapsed = now > submit_ts ? now - submit_ts : 0; - atomic_fetch_add_explicit( - &s->metrics.op_latency.count[kind], 1, memory_order_relaxed); - atomic_fetch_add_explicit( - &s->metrics.op_latency.sum_ns[kind], elapsed, memory_order_relaxed); - atomic_max_u64(&s->metrics.op_latency.max_ns[kind], elapsed); - atomic_fetch_add_explicit( - &s->metrics.op_latency.buckets[kind][latency_bucket(elapsed)], - 1, - memory_order_relaxed); -} - -static int -latency_enabled(const struct damacy_latency_model* l) -{ - return l && (l->baseline_ns || l->lognormal_mu_ln_ns != 0.0 || - l->lognormal_sigma_ln_ns != 0.0); -} - -static uint32_t -pcg32(uint64_t* state) -{ - uint64_t oldstate = *state; - *state = oldstate * 6364136223846793005ULL + 1442695040888963407ULL; - uint32_t xorshifted = (uint32_t)(((oldstate >> 18u) ^ oldstate) >> 27u); - uint32_t rot = (uint32_t)(oldstate >> 59u); - return (xorshifted >> rot) | (xorshifted << ((-rot) & 31u)); -} - -static double -uniform01(uint64_t* state) -{ - uint32_t u = pcg32(state); - return ((double)u + 1.0) / 4294967297.0; -} - -static double -normal01(uint64_t* state) -{ - double u1 = uniform01(state); - double u2 = uniform01(state); - return sqrt(-2.0 * log(u1)) * cos(6.2831853071795864769 * u2); -} - -static uint64_t -sample_delay_ns(struct metadata_store_async* s) -{ - const struct damacy_latency_model* l = &s->latency; - uint64_t ns = l->baseline_ns; - if (l->lognormal_mu_ln_ns == 0.0 && l->lognormal_sigma_ln_ns == 0.0) - return ns; - - pthread_mutex_lock(&s->rng_lock); - double z = normal01(&s->rng_state); - pthread_mutex_unlock(&s->rng_lock); - - double tail = exp(l->lognormal_mu_ln_ns + l->lognormal_sigma_ln_ns * z); - uint64_t tail_ns = tail > 0.0 ? (uint64_t)tail : 0; - if (l->cap_ns && tail_ns > l->cap_ns) - tail_ns = l->cap_ns; - if (UINT64_MAX - ns < tail_ns) - return UINT64_MAX; - return ns + tail_ns; -} - -static void -sleep_for_sample(struct metadata_store_async* s, enum latency_op_kind kind) -{ - if (!s->latency_enabled) - return; - uint64_t ns = sample_delay_ns(s); - if (ns > INT64_MAX) - ns = INT64_MAX; - atomic_fetch_add_explicit(&s->metrics.injector.ops, 1, memory_order_relaxed); - atomic_fetch_add_explicit( - &s->metrics.injector.total_sleep_ns, ns, memory_order_relaxed); - atomic_max_u64(&s->metrics.injector.max_sleep_ns, ns); - switch (kind) { - case LATENCY_OP_STAT: - atomic_fetch_add_explicit( - &s->metrics.injector.stat_ops, 1, memory_order_relaxed); - break; - case LATENCY_OP_SUBMIT: - atomic_fetch_add_explicit( - &s->metrics.injector.submit_ops, 1, memory_order_relaxed); - break; - } - uint64_t active = atomic_fetch_add_explicit( - &s->metrics.injector.active, 1, memory_order_relaxed) + - 1; - atomic_max_u64(&s->metrics.injector.max_active, active); - if (ns) - platform_sleep_ns((int64_t)ns); - atomic_fetch_sub_explicit(&s->metrics.injector.active, 1, memory_order_relaxed); -} - -static void -read_active_begin(struct metadata_store_async* s) -{ - atomic_fetch_add_explicit(&s->metrics.read.jobs, 1, memory_order_relaxed); - uint64_t active = atomic_fetch_add_explicit( - &s->metrics.read.active, 1, memory_order_relaxed) + - 1; - atomic_max_u64(&s->metrics.read.max_active, active); -} - -static void -read_active_end(struct metadata_store_async* s) -{ - atomic_fetch_sub_explicit(&s->metrics.read.active, 1, memory_order_relaxed); -} - -static int -status_not_found_errno(int err) -{ - return err == ENOENT || err == ENOTDIR; -} - -static enum damacy_status -damacy_status_from_errno(int err) -{ - return status_not_found_errno(err) ? DAMACY_NOTFOUND : DAMACY_IO; -} - -static enum store_stat_result -stat_status_from_errno(int err) -{ - return status_not_found_errno(err) ? STORE_STAT_NOT_FOUND : STORE_STAT_ERROR; -} - static void job_free(struct metadata_job* job) { @@ -384,7 +169,7 @@ submit_op(struct metadata_store_async* s, if (!sqe) return 1; set_sqe_data(sqe, job, kind); - job->submit_ts[kind] = monotonic_ns(); + job->submit_ts[kind] = metadata_monotonic_ns(); *sqe_out = sqe; return 0; } @@ -402,7 +187,7 @@ complete_job(struct metadata_store_async* s, struct metadata_job* job) stat_status = STORE_STAT_OK; size = (uint64_t)job->stx.stx_size; } else { - stat_status = stat_status_from_errno(-job->stat_res); + stat_status = metadata_stat_status_from_errno(-job->stat_res); } job->stat_cb(job->user, stat_status, size); } else { @@ -447,7 +232,7 @@ submit_read(struct metadata_store_async* s, struct metadata_job* job) struct io_uring_sqe* sqe = NULL; if (submit_op(s, job, OP_READ, &sqe)) return 1; - read_active_begin(s); + metadata_read_active_begin(&s->common); io_uring_prep_read(sqe, job->fd, job->data, job->len, job->offset); return 0; } @@ -525,13 +310,13 @@ start_job(struct metadata_store_async* s, struct metadata_job* job) job->status = DAMACY_OK; if (job->kind == REQ_STAT) { - sleep_for_sample(s, LATENCY_OP_STAT); + metadata_inject_latency(&s->common, LATENCY_OP_STAT); if (submit_statx(s, job)) fail_request_before_submit(s, job, DAMACY_OOM); return; } - sleep_for_sample(s, LATENCY_OP_SUBMIT); + metadata_inject_latency(&s->common, LATENCY_OP_SUBMIT); if (job->kind == REQ_READ_FILE) { if (ensure_sq_space(s, 2)) { fail_request_before_submit(s, job, DAMACY_IO); @@ -610,7 +395,7 @@ handle_statx_complete(struct metadata_store_async* s, struct metadata_job* job) } if (job->stat_res < 0) { - job->status = damacy_status_from_errno(-job->stat_res); + job->status = metadata_status_from_errno(-job->stat_res); if (job->open_done) { if (job->fd >= 0) finish_job(s, job); @@ -634,7 +419,7 @@ handle_open_complete(struct metadata_store_async* s, struct metadata_job* job) { job->open_done = 1; if (job->open_res < 0) { - job->status = damacy_status_from_errno(-job->open_res); + job->status = metadata_status_from_errno(-job->open_res); if (job->kind == REQ_READ || job->stat_done) complete_job(s, job); return; @@ -657,12 +442,12 @@ handle_open_complete(struct metadata_store_async* s, struct metadata_job* job) static void handle_read_complete(struct metadata_store_async* s, struct metadata_job* job) { - read_active_end(s); + metadata_read_active_end(&s->common); if (job->read_res < 0) { free(job->data); job->data = NULL; job->len = 0; - job->status = damacy_status_from_errno(-job->read_res); + job->status = metadata_status_from_errno(-job->read_res); } else if ((size_t)job->read_res != job->len) { free(job->data); job->data = NULL; @@ -687,7 +472,7 @@ handle_cqe(struct metadata_store_async* s, struct io_uring_cqe* cqe) return; struct metadata_job* job = (struct metadata_job*)(data & ~OP_TAG_MASK); enum op_kind kind = (enum op_kind)(data & OP_TAG_MASK); - record_op_latency(s, kind, job->submit_ts[kind]); + metadata_record_op_latency(&s->common, kind, job->submit_ts[kind]); switch (kind) { case OP_STATX: job->stat_res = cqe->res; @@ -844,8 +629,8 @@ metadata_store_async_free_partial(struct metadata_store_async* s) io_uring_queue_exit(&s->ring); if (s->lock_initialized) pthread_mutex_destroy(&s->lock); - if (s->rng_lock_initialized) - pthread_mutex_destroy(&s->rng_lock); + if (s->common_initialized) + metadata_store_common_destroy(&s->common); while (s->pending_head) { struct metadata_job* job = s->pending_head; s->pending_head = job->next; @@ -881,18 +666,13 @@ metadata_store_async_create(int concurrency, #else (void)affinity; #endif - if (latency) { - s->latency = *latency; - s->latency_enabled = latency_enabled(latency); - s->rng_state = latency->seed ? latency->seed : 0xc0ffee1234ULL; - } if (pthread_mutex_init(&s->lock, NULL) != 0) goto Fail; s->lock_initialized = 1; - if (pthread_mutex_init(&s->rng_lock, NULL) != 0) + if (metadata_store_common_init(&s->common, latency)) goto Fail; - s->rng_lock_initialized = 1; + s->common_initialized = 1; uint32_t entries = ring_entries_for_concurrency(concurrency); int rc = io_uring_queue_init(entries, &s->ring, 0); @@ -967,6 +747,12 @@ metadata_store_async_destroy(struct metadata_store_async* s) metadata_store_async_free_partial(s); } +struct metadata_store_common* +metadata_store_async_common(struct metadata_store_async* s) +{ + return &s->common; +} + static int enqueue_job(struct metadata_store_async* s, struct metadata_job* job) { @@ -1050,122 +836,3 @@ metadata_store_async_stat(struct metadata_store_async* s, } return 0; } - -void -metadata_store_async_latency_stats_get( - struct metadata_store_async* s, - struct metadata_store_async_latency_stats* out) -{ - if (!out) - return; - *out = (struct metadata_store_async_latency_stats){ 0 }; - if (!s) - return; - *out = (struct metadata_store_async_latency_stats){ - .ops = atomic_load_explicit(&s->metrics.injector.ops, memory_order_relaxed), - .stat_ops = - atomic_load_explicit(&s->metrics.injector.stat_ops, memory_order_relaxed), - .submit_ops = atomic_load_explicit(&s->metrics.injector.submit_ops, - memory_order_relaxed), - .active = - atomic_load_explicit(&s->metrics.injector.active, memory_order_relaxed), - .max_active = atomic_load_explicit(&s->metrics.injector.max_active, - memory_order_relaxed), - .total_sleep_ns = atomic_load_explicit(&s->metrics.injector.total_sleep_ns, - memory_order_relaxed), - .max_sleep_ns = atomic_load_explicit(&s->metrics.injector.max_sleep_ns, - memory_order_relaxed), - }; -} - -void -metadata_store_async_latency_stats_reset(struct metadata_store_async* s) -{ - if (!s) - return; - uint64_t active = - atomic_load_explicit(&s->metrics.injector.active, memory_order_relaxed); - atomic_store_explicit(&s->metrics.injector.ops, 0, memory_order_relaxed); - atomic_store_explicit(&s->metrics.injector.stat_ops, 0, memory_order_relaxed); - atomic_store_explicit( - &s->metrics.injector.submit_ops, 0, memory_order_relaxed); - atomic_store_explicit( - &s->metrics.injector.max_active, active, memory_order_relaxed); - atomic_store_explicit( - &s->metrics.injector.total_sleep_ns, 0, memory_order_relaxed); - atomic_store_explicit( - &s->metrics.injector.max_sleep_ns, 0, memory_order_relaxed); -} - -void -metadata_store_async_backend_stats_get( - struct metadata_store_async* s, - struct metadata_store_async_backend_stats* out) -{ - if (!out) - return; - *out = (struct metadata_store_async_backend_stats){ 0 }; - if (!s) - return; - *out = (struct metadata_store_async_backend_stats){ - .read_jobs = - atomic_load_explicit(&s->metrics.read.jobs, memory_order_relaxed), - .read_active = - atomic_load_explicit(&s->metrics.read.active, memory_order_relaxed), - .read_max_active = - atomic_load_explicit(&s->metrics.read.max_active, memory_order_relaxed), - }; -} - -void -metadata_store_async_backend_stats_reset(struct metadata_store_async* s) -{ - if (!s) - return; - uint64_t active = - atomic_load_explicit(&s->metrics.read.active, memory_order_relaxed); - atomic_store_explicit(&s->metrics.read.jobs, 0, memory_order_relaxed); - atomic_store_explicit( - &s->metrics.read.max_active, active, memory_order_relaxed); -} - -void -metadata_store_async_op_latency_stats_get( - struct metadata_store_async* s, - struct metadata_store_async_op_latency_stats* out) -{ - if (!out) - return; - *out = (struct metadata_store_async_op_latency_stats){ 0 }; - if (!s) - return; - for (unsigned k = 0; k < METADATA_OP_LATENCY_NKINDS; ++k) { - out->kinds[k].count = atomic_load_explicit(&s->metrics.op_latency.count[k], - memory_order_relaxed); - out->kinds[k].sum_ns = atomic_load_explicit( - &s->metrics.op_latency.sum_ns[k], memory_order_relaxed); - out->kinds[k].max_ns = atomic_load_explicit( - &s->metrics.op_latency.max_ns[k], memory_order_relaxed); - for (unsigned b = 0; b < METADATA_OP_LATENCY_NBUCKETS; ++b) - out->kinds[k].buckets[b] = atomic_load_explicit( - &s->metrics.op_latency.buckets[k][b], memory_order_relaxed); - } -} - -void -metadata_store_async_op_latency_stats_reset(struct metadata_store_async* s) -{ - if (!s) - return; - for (unsigned k = 0; k < METADATA_OP_LATENCY_NKINDS; ++k) { - atomic_store_explicit( - &s->metrics.op_latency.count[k], 0, memory_order_relaxed); - atomic_store_explicit( - &s->metrics.op_latency.sum_ns[k], 0, memory_order_relaxed); - atomic_store_explicit( - &s->metrics.op_latency.max_ns[k], 0, memory_order_relaxed); - for (unsigned b = 0; b < METADATA_OP_LATENCY_NBUCKETS; ++b) - atomic_store_explicit( - &s->metrics.op_latency.buckets[k][b], 0, memory_order_relaxed); - } -} diff --git a/src/store/metadata_store_async.posix.c b/src/store/metadata_store_async.posix.c index 3c0f7767..49b68fb2 100644 --- a/src/store/metadata_store_async.posix.c +++ b/src/store/metadata_store_async.posix.c @@ -2,16 +2,14 @@ // Bulk data still uses the unchanged store_fs/io_queue path. #include "store/metadata_store_async.h" #include "log/log.h" -#include "platform/platform.h" +#include "store/metadata_store_async_common.h" #include #include #include #include -#include #include #include #include -#include #include enum request_kind @@ -21,50 +19,6 @@ enum request_kind REQ_STAT, }; -enum op_kind -{ - OP_STATX, - OP_OPEN, - OP_READ, - OP_CLOSE, -}; - -enum latency_op_kind -{ - LATENCY_OP_STAT, - LATENCY_OP_SUBMIT, -}; - -struct metadata_metrics -{ - struct - { - _Atomic uint64_t ops; - _Atomic uint64_t stat_ops; - _Atomic uint64_t submit_ops; - _Atomic uint64_t active; - _Atomic uint64_t max_active; - _Atomic uint64_t total_sleep_ns; - _Atomic uint64_t max_sleep_ns; - } injector; - - struct - { - _Atomic uint64_t jobs; - _Atomic uint64_t active; - _Atomic uint64_t max_active; - } read; - - struct - { - _Atomic uint64_t count[METADATA_OP_LATENCY_NKINDS]; - _Atomic uint64_t sum_ns[METADATA_OP_LATENCY_NKINDS]; - _Atomic uint64_t max_ns[METADATA_OP_LATENCY_NKINDS]; - _Atomic uint64_t buckets[METADATA_OP_LATENCY_NKINDS] - [METADATA_OP_LATENCY_NBUCKETS]; - } op_latency; -}; - struct metadata_job { struct metadata_job* next; @@ -84,181 +38,9 @@ struct metadata_store_async pthread_cond_t wake; struct metadata_job *pending_head, *pending_tail; int stopping; - struct damacy_latency_model latency; - int latency_enabled; - uint64_t rng_state; - pthread_mutex_t rng_lock; - struct metadata_metrics metrics; + struct metadata_store_common common; }; -static void -atomic_max_u64(_Atomic uint64_t* dst, uint64_t val) -{ - uint64_t cur = atomic_load_explicit(dst, memory_order_relaxed); - while (cur < val && - !atomic_compare_exchange_weak_explicit( - dst, &cur, val, memory_order_relaxed, memory_order_relaxed)) { - } -} - -static uint64_t -monotonic_ns(void) -{ - struct timespec now; - clock_gettime(CLOCK_MONOTONIC, &now); - return (uint64_t)now.tv_sec * 1000000000ull + (uint64_t)now.tv_nsec; -} - -static unsigned -latency_bucket(uint64_t ns) -{ - if (ns == 0) - return 0; - unsigned idx = 63u - (unsigned)__builtin_clzll(ns); - if (idx >= METADATA_OP_LATENCY_NBUCKETS) - idx = METADATA_OP_LATENCY_NBUCKETS - 1; - return idx; -} - -static void -record_op_latency(struct metadata_store_async* s, - enum op_kind kind, - uint64_t submit_ts) -{ - if (!submit_ts) - return; - uint64_t now = monotonic_ns(); - uint64_t elapsed = now > submit_ts ? now - submit_ts : 0; - atomic_fetch_add_explicit( - &s->metrics.op_latency.count[kind], 1, memory_order_relaxed); - atomic_fetch_add_explicit( - &s->metrics.op_latency.sum_ns[kind], elapsed, memory_order_relaxed); - atomic_max_u64(&s->metrics.op_latency.max_ns[kind], elapsed); - atomic_fetch_add_explicit( - &s->metrics.op_latency.buckets[kind][latency_bucket(elapsed)], - 1, - memory_order_relaxed); -} - -static int -latency_enabled(const struct damacy_latency_model* l) -{ - return l && (l->baseline_ns || l->lognormal_mu_ln_ns != 0.0 || - l->lognormal_sigma_ln_ns != 0.0); -} - -static uint32_t -pcg32(uint64_t* state) -{ - uint64_t oldstate = *state; - *state = oldstate * 6364136223846793005ULL + 1442695040888963407ULL; - uint32_t xorshifted = (uint32_t)(((oldstate >> 18u) ^ oldstate) >> 27u); - uint32_t rot = (uint32_t)(oldstate >> 59u); - return (xorshifted >> rot) | (xorshifted << ((-rot) & 31u)); -} - -static double -uniform01(uint64_t* state) -{ - uint32_t u = pcg32(state); - return ((double)u + 1.0) / 4294967297.0; -} - -static double -normal01(uint64_t* state) -{ - double u1 = uniform01(state); - double u2 = uniform01(state); - return sqrt(-2.0 * log(u1)) * cos(6.2831853071795864769 * u2); -} - -static uint64_t -sample_delay_ns(struct metadata_store_async* s) -{ - const struct damacy_latency_model* l = &s->latency; - uint64_t ns = l->baseline_ns; - if (l->lognormal_mu_ln_ns == 0.0 && l->lognormal_sigma_ln_ns == 0.0) - return ns; - - pthread_mutex_lock(&s->rng_lock); - double z = normal01(&s->rng_state); - pthread_mutex_unlock(&s->rng_lock); - - double tail = exp(l->lognormal_mu_ln_ns + l->lognormal_sigma_ln_ns * z); - uint64_t tail_ns = tail > 0.0 ? (uint64_t)tail : 0; - if (l->cap_ns && tail_ns > l->cap_ns) - tail_ns = l->cap_ns; - if (UINT64_MAX - ns < tail_ns) - return UINT64_MAX; - return ns + tail_ns; -} - -static void -sleep_for_sample(struct metadata_store_async* s, enum latency_op_kind kind) -{ - if (!s->latency_enabled) - return; - uint64_t ns = sample_delay_ns(s); - if (ns > INT64_MAX) - ns = INT64_MAX; - atomic_fetch_add_explicit(&s->metrics.injector.ops, 1, memory_order_relaxed); - atomic_fetch_add_explicit( - &s->metrics.injector.total_sleep_ns, ns, memory_order_relaxed); - atomic_max_u64(&s->metrics.injector.max_sleep_ns, ns); - switch (kind) { - case LATENCY_OP_STAT: - atomic_fetch_add_explicit( - &s->metrics.injector.stat_ops, 1, memory_order_relaxed); - break; - case LATENCY_OP_SUBMIT: - atomic_fetch_add_explicit( - &s->metrics.injector.submit_ops, 1, memory_order_relaxed); - break; - } - uint64_t active = atomic_fetch_add_explicit( - &s->metrics.injector.active, 1, memory_order_relaxed) + - 1; - atomic_max_u64(&s->metrics.injector.max_active, active); - if (ns) - platform_sleep_ns((int64_t)ns); - atomic_fetch_sub_explicit( - &s->metrics.injector.active, 1, memory_order_relaxed); -} - -static void -read_active_begin(struct metadata_store_async* s) -{ - atomic_fetch_add_explicit(&s->metrics.read.jobs, 1, memory_order_relaxed); - uint64_t active = atomic_fetch_add_explicit( - &s->metrics.read.active, 1, memory_order_relaxed) + - 1; - atomic_max_u64(&s->metrics.read.max_active, active); -} - -static void -read_active_end(struct metadata_store_async* s) -{ - atomic_fetch_sub_explicit(&s->metrics.read.active, 1, memory_order_relaxed); -} - -static int -status_not_found_errno(int err) -{ - return err == ENOENT || err == ENOTDIR; -} - -static enum damacy_status -damacy_status_from_errno(int err) -{ - return status_not_found_errno(err) ? DAMACY_NOTFOUND : DAMACY_IO; -} - -static enum store_stat_result -stat_status_from_errno(int err) -{ - return status_not_found_errno(err) ? STORE_STAT_NOT_FOUND : STORE_STAT_ERROR; -} - static void job_free(struct metadata_job* job) { @@ -284,33 +66,33 @@ job_new(const char* key, enum request_kind kind) static void process_job(struct metadata_store_async* s, struct metadata_job* job) { - sleep_for_sample(s, - job->kind == REQ_STAT ? LATENCY_OP_STAT : LATENCY_OP_SUBMIT); + metadata_inject_latency( + &s->common, job->kind == REQ_STAT ? LATENCY_OP_STAT : LATENCY_OP_SUBMIT); if (job->kind == REQ_STAT) { struct stat st; - uint64_t started = monotonic_ns(); + uint64_t started = metadata_monotonic_ns(); int rc = stat(job->key, &st), error = errno; - record_op_latency(s, OP_STATX, started); + metadata_record_op_latency(&s->common, OP_STATX, started); job->stat_cb(job->user, - rc ? stat_status_from_errno(error) : STORE_STAT_OK, + rc ? metadata_stat_status_from_errno(error) : STORE_STAT_OK, rc ? 0 : (uint64_t)st.st_size); return; } enum damacy_status status = DAMACY_OK; void* data = NULL; size_t len = job->requested_len; - uint64_t started = monotonic_ns(); + uint64_t started = metadata_monotonic_ns(); int fd = open(job->key, O_RDONLY | O_CLOEXEC), error = errno; - record_op_latency(s, OP_OPEN, started); + metadata_record_op_latency(&s->common, OP_OPEN, started); if (fd < 0) { - status = damacy_status_from_errno(error); + status = metadata_status_from_errno(error); goto Complete; } if (job->kind == REQ_READ_FILE) { struct stat st; - started = monotonic_ns(); + started = metadata_monotonic_ns(); int rc = fstat(fd, &st); - record_op_latency(s, OP_STATX, started); + metadata_record_op_latency(&s->common, OP_STATX, started); if (rc || st.st_size < 0 || (uint64_t)st.st_size > SIZE_MAX) { status = DAMACY_IO; goto Close; @@ -328,8 +110,8 @@ process_job(struct metadata_store_async* s, struct metadata_job* job) status = DAMACY_OOM; goto Close; } - read_active_begin(s); - started = monotonic_ns(); + metadata_read_active_begin(&s->common); + started = metadata_monotonic_ns(); size_t done = 0; while (done < len) { ssize_t n = @@ -342,15 +124,15 @@ process_job(struct metadata_store_async* s, struct metadata_job* job) } done += (size_t)n; } - record_op_latency(s, OP_READ, started); - read_active_end(s); + metadata_record_op_latency(&s->common, OP_READ, started); + metadata_read_active_end(&s->common); Close: - started = monotonic_ns(); + started = metadata_monotonic_ns(); if (close(fd) && status == DAMACY_OK) log_warn("metadata_store_async: close failed for %s: %s", job->key, strerror(errno)); - record_op_latency(s, OP_CLOSE, started); + metadata_record_op_latency(&s->common, OP_CLOSE, started); Complete: if (status != DAMACY_OK) { free(data); @@ -394,13 +176,19 @@ metadata_store_async_destroy(struct metadata_store_async* s) pthread_mutex_unlock(&s->lock); for (int i = 0; i < s->started; ++i) pthread_join(s->workers[i], NULL); - pthread_mutex_destroy(&s->rng_lock); + metadata_store_common_destroy(&s->common); pthread_cond_destroy(&s->wake); pthread_mutex_destroy(&s->lock); free(s->workers); free(s); } +struct metadata_store_common* +metadata_store_async_common(struct metadata_store_async* s) +{ + return &s->common; +} + struct metadata_store_async* metadata_store_async_create(int concurrency, const struct numa_resolved* affinity, @@ -427,15 +215,11 @@ metadata_store_async_create(int concurrency, pthread_mutex_destroy(&s->lock); goto Fail; } - if (pthread_mutex_init(&s->rng_lock, NULL)) { + if (metadata_store_common_init(&s->common, latency)) { pthread_cond_destroy(&s->wake); pthread_mutex_destroy(&s->lock); goto Fail; } - if (latency) - s->latency = *latency; - s->latency_enabled = latency_enabled(latency); - s->rng_state = s->latency.seed ? s->latency.seed : 0xc0ffee1234ULL; for (int i = 0; i < concurrency; ++i) { if (pthread_create(&s->workers[i], NULL, worker_main, s)) { metadata_store_async_destroy(s); @@ -532,122 +316,3 @@ metadata_store_async_stat(struct metadata_store_async* s, } return 0; } - -void -metadata_store_async_latency_stats_get( - struct metadata_store_async* s, - struct metadata_store_async_latency_stats* out) -{ - if (!out) - return; - *out = (struct metadata_store_async_latency_stats){ 0 }; - if (!s) - return; - *out = (struct metadata_store_async_latency_stats){ - .ops = atomic_load_explicit(&s->metrics.injector.ops, memory_order_relaxed), - .stat_ops = - atomic_load_explicit(&s->metrics.injector.stat_ops, memory_order_relaxed), - .submit_ops = atomic_load_explicit(&s->metrics.injector.submit_ops, - memory_order_relaxed), - .active = - atomic_load_explicit(&s->metrics.injector.active, memory_order_relaxed), - .max_active = atomic_load_explicit(&s->metrics.injector.max_active, - memory_order_relaxed), - .total_sleep_ns = atomic_load_explicit(&s->metrics.injector.total_sleep_ns, - memory_order_relaxed), - .max_sleep_ns = atomic_load_explicit(&s->metrics.injector.max_sleep_ns, - memory_order_relaxed), - }; -} - -void -metadata_store_async_latency_stats_reset(struct metadata_store_async* s) -{ - if (!s) - return; - uint64_t active = - atomic_load_explicit(&s->metrics.injector.active, memory_order_relaxed); - atomic_store_explicit(&s->metrics.injector.ops, 0, memory_order_relaxed); - atomic_store_explicit(&s->metrics.injector.stat_ops, 0, memory_order_relaxed); - atomic_store_explicit( - &s->metrics.injector.submit_ops, 0, memory_order_relaxed); - atomic_store_explicit( - &s->metrics.injector.max_active, active, memory_order_relaxed); - atomic_store_explicit( - &s->metrics.injector.total_sleep_ns, 0, memory_order_relaxed); - atomic_store_explicit( - &s->metrics.injector.max_sleep_ns, 0, memory_order_relaxed); -} - -void -metadata_store_async_backend_stats_get( - struct metadata_store_async* s, - struct metadata_store_async_backend_stats* out) -{ - if (!out) - return; - *out = (struct metadata_store_async_backend_stats){ 0 }; - if (!s) - return; - *out = (struct metadata_store_async_backend_stats){ - .read_jobs = - atomic_load_explicit(&s->metrics.read.jobs, memory_order_relaxed), - .read_active = - atomic_load_explicit(&s->metrics.read.active, memory_order_relaxed), - .read_max_active = - atomic_load_explicit(&s->metrics.read.max_active, memory_order_relaxed), - }; -} - -void -metadata_store_async_backend_stats_reset(struct metadata_store_async* s) -{ - if (!s) - return; - uint64_t active = - atomic_load_explicit(&s->metrics.read.active, memory_order_relaxed); - atomic_store_explicit(&s->metrics.read.jobs, 0, memory_order_relaxed); - atomic_store_explicit( - &s->metrics.read.max_active, active, memory_order_relaxed); -} - -void -metadata_store_async_op_latency_stats_get( - struct metadata_store_async* s, - struct metadata_store_async_op_latency_stats* out) -{ - if (!out) - return; - *out = (struct metadata_store_async_op_latency_stats){ 0 }; - if (!s) - return; - for (unsigned k = 0; k < METADATA_OP_LATENCY_NKINDS; ++k) { - out->kinds[k].count = atomic_load_explicit(&s->metrics.op_latency.count[k], - memory_order_relaxed); - out->kinds[k].sum_ns = atomic_load_explicit( - &s->metrics.op_latency.sum_ns[k], memory_order_relaxed); - out->kinds[k].max_ns = atomic_load_explicit( - &s->metrics.op_latency.max_ns[k], memory_order_relaxed); - for (unsigned b = 0; b < METADATA_OP_LATENCY_NBUCKETS; ++b) - out->kinds[k].buckets[b] = atomic_load_explicit( - &s->metrics.op_latency.buckets[k][b], memory_order_relaxed); - } -} - -void -metadata_store_async_op_latency_stats_reset(struct metadata_store_async* s) -{ - if (!s) - return; - for (unsigned k = 0; k < METADATA_OP_LATENCY_NKINDS; ++k) { - atomic_store_explicit( - &s->metrics.op_latency.count[k], 0, memory_order_relaxed); - atomic_store_explicit( - &s->metrics.op_latency.sum_ns[k], 0, memory_order_relaxed); - atomic_store_explicit( - &s->metrics.op_latency.max_ns[k], 0, memory_order_relaxed); - for (unsigned b = 0; b < METADATA_OP_LATENCY_NBUCKETS; ++b) - atomic_store_explicit( - &s->metrics.op_latency.buckets[k][b], 0, memory_order_relaxed); - } -} diff --git a/src/store/metadata_store_async_common.c b/src/store/metadata_store_async_common.c new file mode 100644 index 00000000..ffe66233 --- /dev/null +++ b/src/store/metadata_store_async_common.c @@ -0,0 +1,306 @@ +#include "store/metadata_store_async_common.h" + +#include "platform/platform.h" + +#include +#include +#include + +static void +atomic_max_u64(_Atomic uint64_t* dst, uint64_t val) +{ + uint64_t cur = atomic_load_explicit(dst, memory_order_relaxed); + while (cur < val && + !atomic_compare_exchange_weak_explicit( + dst, &cur, val, memory_order_relaxed, memory_order_relaxed)) { + } +} + +uint64_t +metadata_monotonic_ns(void) +{ + struct timespec now; + clock_gettime(CLOCK_MONOTONIC, &now); + return (uint64_t)now.tv_sec * 1000000000ull + (uint64_t)now.tv_nsec; +} + +static unsigned +latency_bucket(uint64_t ns) +{ + if (ns == 0) + return 0; + unsigned idx = 63u - (unsigned)__builtin_clzll(ns); + if (idx >= METADATA_OP_LATENCY_NBUCKETS) + idx = METADATA_OP_LATENCY_NBUCKETS - 1; + return idx; +} + +void +metadata_record_op_latency(struct metadata_store_common* c, + enum op_kind kind, + uint64_t submit_ts) +{ + if (!submit_ts) + return; + uint64_t now = metadata_monotonic_ns(); + uint64_t elapsed = now > submit_ts ? now - submit_ts : 0; + atomic_fetch_add_explicit( + &c->metrics.op_latency.count[kind], 1, memory_order_relaxed); + atomic_fetch_add_explicit( + &c->metrics.op_latency.sum_ns[kind], elapsed, memory_order_relaxed); + atomic_max_u64(&c->metrics.op_latency.max_ns[kind], elapsed); + atomic_fetch_add_explicit( + &c->metrics.op_latency.buckets[kind][latency_bucket(elapsed)], + 1, + memory_order_relaxed); +} + +static int +latency_enabled(const struct damacy_latency_model* l) +{ + return l && (l->baseline_ns || l->lognormal_mu_ln_ns != 0.0 || + l->lognormal_sigma_ln_ns != 0.0); +} + +int +metadata_store_common_init(struct metadata_store_common* c, + const struct damacy_latency_model* latency) +{ + if (latency) + c->latency = *latency; + c->latency_enabled = latency_enabled(latency); + c->rng_state = c->latency.seed ? c->latency.seed : 0xc0ffee1234ULL; + return pthread_mutex_init(&c->rng_lock, NULL) != 0; +} + +void +metadata_store_common_destroy(struct metadata_store_common* c) +{ + pthread_mutex_destroy(&c->rng_lock); +} + +static uint32_t +pcg32(uint64_t* state) +{ + uint64_t oldstate = *state; + *state = oldstate * 6364136223846793005ULL + 1442695040888963407ULL; + uint32_t xorshifted = (uint32_t)(((oldstate >> 18u) ^ oldstate) >> 27u); + uint32_t rot = (uint32_t)(oldstate >> 59u); + return (xorshifted >> rot) | (xorshifted << ((-rot) & 31u)); +} + +static double +uniform01(uint64_t* state) +{ + uint32_t u = pcg32(state); + return ((double)u + 1.0) / 4294967297.0; +} + +static double +normal01(uint64_t* state) +{ + double u1 = uniform01(state); + double u2 = uniform01(state); + return sqrt(-2.0 * log(u1)) * cos(6.2831853071795864769 * u2); +} + +static uint64_t +sample_delay_ns(struct metadata_store_common* c) +{ + const struct damacy_latency_model* l = &c->latency; + uint64_t ns = l->baseline_ns; + if (l->lognormal_mu_ln_ns == 0.0 && l->lognormal_sigma_ln_ns == 0.0) + return ns; + + pthread_mutex_lock(&c->rng_lock); + double z = normal01(&c->rng_state); + pthread_mutex_unlock(&c->rng_lock); + + double tail = exp(l->lognormal_mu_ln_ns + l->lognormal_sigma_ln_ns * z); + uint64_t tail_ns = tail > 0.0 ? (uint64_t)tail : 0; + if (l->cap_ns && tail_ns > l->cap_ns) + tail_ns = l->cap_ns; + if (UINT64_MAX - ns < tail_ns) + return UINT64_MAX; + return ns + tail_ns; +} + +void +metadata_inject_latency(struct metadata_store_common* c, + enum latency_op_kind kind) +{ + if (!c->latency_enabled) + return; + uint64_t ns = sample_delay_ns(c); + if (ns > INT64_MAX) + ns = INT64_MAX; + atomic_fetch_add_explicit(&c->metrics.injector.ops, 1, memory_order_relaxed); + atomic_fetch_add_explicit( + &c->metrics.injector.total_sleep_ns, ns, memory_order_relaxed); + atomic_max_u64(&c->metrics.injector.max_sleep_ns, ns); + switch (kind) { + case LATENCY_OP_STAT: + atomic_fetch_add_explicit( + &c->metrics.injector.stat_ops, 1, memory_order_relaxed); + break; + case LATENCY_OP_SUBMIT: + atomic_fetch_add_explicit( + &c->metrics.injector.submit_ops, 1, memory_order_relaxed); + break; + } + uint64_t active = atomic_fetch_add_explicit( + &c->metrics.injector.active, 1, memory_order_relaxed) + + 1; + atomic_max_u64(&c->metrics.injector.max_active, active); + if (ns) + platform_sleep_ns((int64_t)ns); + atomic_fetch_sub_explicit( + &c->metrics.injector.active, 1, memory_order_relaxed); +} + +void +metadata_read_active_begin(struct metadata_store_common* c) +{ + atomic_fetch_add_explicit(&c->metrics.read.jobs, 1, memory_order_relaxed); + uint64_t active = atomic_fetch_add_explicit( + &c->metrics.read.active, 1, memory_order_relaxed) + + 1; + atomic_max_u64(&c->metrics.read.max_active, active); +} + +void +metadata_read_active_end(struct metadata_store_common* c) +{ + atomic_fetch_sub_explicit(&c->metrics.read.active, 1, memory_order_relaxed); +} + +static int +status_not_found_errno(int err) +{ + return err == ENOENT || err == ENOTDIR; +} + +enum damacy_status +metadata_status_from_errno(int err) +{ + return status_not_found_errno(err) ? DAMACY_NOTFOUND : DAMACY_IO; +} + +enum store_stat_result +metadata_stat_status_from_errno(int err) +{ + return status_not_found_errno(err) ? STORE_STAT_NOT_FOUND : STORE_STAT_ERROR; +} + +void +metadata_store_async_latency_stats_get( + struct metadata_store_async* s, + struct metadata_store_async_latency_stats* out) +{ + if (!out) + return; + *out = (struct metadata_store_async_latency_stats){ 0 }; + if (!s) + return; + struct metadata_metrics* m = &metadata_store_async_common(s)->metrics; + *out = (struct metadata_store_async_latency_stats){ + .ops = atomic_load_explicit(&m->injector.ops, memory_order_relaxed), + .stat_ops = + atomic_load_explicit(&m->injector.stat_ops, memory_order_relaxed), + .submit_ops = + atomic_load_explicit(&m->injector.submit_ops, memory_order_relaxed), + .active = atomic_load_explicit(&m->injector.active, memory_order_relaxed), + .max_active = + atomic_load_explicit(&m->injector.max_active, memory_order_relaxed), + .total_sleep_ns = + atomic_load_explicit(&m->injector.total_sleep_ns, memory_order_relaxed), + .max_sleep_ns = + atomic_load_explicit(&m->injector.max_sleep_ns, memory_order_relaxed), + }; +} + +void +metadata_store_async_latency_stats_reset(struct metadata_store_async* s) +{ + if (!s) + return; + struct metadata_metrics* m = &metadata_store_async_common(s)->metrics; + uint64_t active = + atomic_load_explicit(&m->injector.active, memory_order_relaxed); + atomic_store_explicit(&m->injector.ops, 0, memory_order_relaxed); + atomic_store_explicit(&m->injector.stat_ops, 0, memory_order_relaxed); + atomic_store_explicit(&m->injector.submit_ops, 0, memory_order_relaxed); + atomic_store_explicit(&m->injector.max_active, active, memory_order_relaxed); + atomic_store_explicit(&m->injector.total_sleep_ns, 0, memory_order_relaxed); + atomic_store_explicit(&m->injector.max_sleep_ns, 0, memory_order_relaxed); +} + +void +metadata_store_async_backend_stats_get( + struct metadata_store_async* s, + struct metadata_store_async_backend_stats* out) +{ + if (!out) + return; + *out = (struct metadata_store_async_backend_stats){ 0 }; + if (!s) + return; + struct metadata_metrics* m = &metadata_store_async_common(s)->metrics; + *out = (struct metadata_store_async_backend_stats){ + .read_jobs = atomic_load_explicit(&m->read.jobs, memory_order_relaxed), + .read_active = atomic_load_explicit(&m->read.active, memory_order_relaxed), + .read_max_active = + atomic_load_explicit(&m->read.max_active, memory_order_relaxed), + }; +} + +void +metadata_store_async_backend_stats_reset(struct metadata_store_async* s) +{ + if (!s) + return; + struct metadata_metrics* m = &metadata_store_async_common(s)->metrics; + uint64_t active = atomic_load_explicit(&m->read.active, memory_order_relaxed); + atomic_store_explicit(&m->read.jobs, 0, memory_order_relaxed); + atomic_store_explicit(&m->read.max_active, active, memory_order_relaxed); +} + +void +metadata_store_async_op_latency_stats_get( + struct metadata_store_async* s, + struct metadata_store_async_op_latency_stats* out) +{ + if (!out) + return; + *out = (struct metadata_store_async_op_latency_stats){ 0 }; + if (!s) + return; + struct metadata_metrics* m = &metadata_store_async_common(s)->metrics; + for (unsigned k = 0; k < METADATA_OP_LATENCY_NKINDS; ++k) { + out->kinds[k].count = + atomic_load_explicit(&m->op_latency.count[k], memory_order_relaxed); + out->kinds[k].sum_ns = + atomic_load_explicit(&m->op_latency.sum_ns[k], memory_order_relaxed); + out->kinds[k].max_ns = + atomic_load_explicit(&m->op_latency.max_ns[k], memory_order_relaxed); + for (unsigned b = 0; b < METADATA_OP_LATENCY_NBUCKETS; ++b) + out->kinds[k].buckets[b] = atomic_load_explicit( + &m->op_latency.buckets[k][b], memory_order_relaxed); + } +} + +void +metadata_store_async_op_latency_stats_reset(struct metadata_store_async* s) +{ + if (!s) + return; + struct metadata_metrics* m = &metadata_store_async_common(s)->metrics; + for (unsigned k = 0; k < METADATA_OP_LATENCY_NKINDS; ++k) { + atomic_store_explicit(&m->op_latency.count[k], 0, memory_order_relaxed); + atomic_store_explicit(&m->op_latency.sum_ns[k], 0, memory_order_relaxed); + atomic_store_explicit(&m->op_latency.max_ns[k], 0, memory_order_relaxed); + for (unsigned b = 0; b < METADATA_OP_LATENCY_NBUCKETS; ++b) + atomic_store_explicit( + &m->op_latency.buckets[k][b], 0, memory_order_relaxed); + } +} diff --git a/src/store/metadata_store_async_common.h b/src/store/metadata_store_async_common.h new file mode 100644 index 00000000..d339b058 --- /dev/null +++ b/src/store/metadata_store_async_common.h @@ -0,0 +1,91 @@ +// Backend-author header: metrics, injected latency, and error mapping shared +// by the io_uring and POSIX metadata backends. +#pragma once + +#include "store/metadata_store_async.h" + +#include +#include +#include + +enum op_kind +{ + OP_STATX, + OP_OPEN, + OP_READ, + OP_CLOSE, +}; + +enum latency_op_kind +{ + LATENCY_OP_STAT, + LATENCY_OP_SUBMIT, +}; + +struct metadata_metrics +{ + struct + { + _Atomic uint64_t ops; + _Atomic uint64_t stat_ops; + _Atomic uint64_t submit_ops; + _Atomic uint64_t active; + _Atomic uint64_t max_active; + _Atomic uint64_t total_sleep_ns; + _Atomic uint64_t max_sleep_ns; + } injector; + + struct + { + _Atomic uint64_t jobs; + _Atomic uint64_t active; + _Atomic uint64_t max_active; + } read; + + struct + { + _Atomic uint64_t count[METADATA_OP_LATENCY_NKINDS]; + _Atomic uint64_t sum_ns[METADATA_OP_LATENCY_NKINDS]; + _Atomic uint64_t max_ns[METADATA_OP_LATENCY_NKINDS]; + _Atomic uint64_t buckets[METADATA_OP_LATENCY_NKINDS] + [METADATA_OP_LATENCY_NBUCKETS]; + } op_latency; +}; + +struct metadata_store_common +{ + struct damacy_latency_model latency; + int latency_enabled; + uint64_t rng_state; + pthread_mutex_t rng_lock; + struct metadata_metrics metrics; +}; + +// Each backend defines this; the stats functions reach the metrics through it. +struct metadata_store_common* +metadata_store_async_common(struct metadata_store_async* s); + +int +metadata_store_common_init(struct metadata_store_common* c, + const struct damacy_latency_model* latency); +void +metadata_store_common_destroy(struct metadata_store_common* c); + +uint64_t +metadata_monotonic_ns(void); +void +metadata_record_op_latency(struct metadata_store_common* c, + enum op_kind kind, + uint64_t submit_ts); +void +metadata_inject_latency(struct metadata_store_common* c, + enum latency_op_kind kind); +void +metadata_read_active_begin(struct metadata_store_common* c); +void +metadata_read_active_end(struct metadata_store_common* c); + +enum damacy_status +metadata_status_from_errno(int err); +enum store_stat_result +metadata_stat_status_from_errno(int err); From 7b261af4933346cf1a0c8691194588089c302c40 Mon Sep 17 00:00:00 2001 From: Nathan Clack Date: Thu, 24 Sep 2026 21:21:08 +0000 Subject: [PATCH 7/9] store: log macOS metadata setup --- src/store/metadata_store_async.posix.c | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/src/store/metadata_store_async.posix.c b/src/store/metadata_store_async.posix.c index 49b68fb2..8839869c 100644 --- a/src/store/metadata_store_async.posix.c +++ b/src/store/metadata_store_async.posix.c @@ -221,12 +221,21 @@ metadata_store_async_create(int concurrency, goto Fail; } for (int i = 0; i < concurrency; ++i) { - if (pthread_create(&s->workers[i], NULL, worker_main, s)) { + int rc = pthread_create(&s->workers[i], NULL, worker_main, s); + if (rc) { + log_error("metadata_store_async: pthread_create failed for worker %d " + "of %d: %s", + i + 1, + concurrency, + strerror(rc)); metadata_store_async_destroy(s); return NULL; } ++s->started; } + log_info("metadata_store_async: using POSIX worker metadata path " + "(workers=%d)", + concurrency); return s; Fail: free(s->workers); From 6af90371014c482b18b250f647c499f1d977bf0e Mon Sep 17 00:00:00 2001 From: Nathan Clack Date: Thu, 24 Sep 2026 21:23:08 +0000 Subject: [PATCH 8/9] reader: reject more workers than CPUs --- docs/pipeline.md | 4 +++- python/damacy/__init__.py | 3 ++- src/damacy.h | 6 +++--- src/pipeline/components.c | 9 +++++++++ tests/test_cpu_pipeline.c | 13 +++++++++++++ 5 files changed, 30 insertions(+), 5 deletions(-) diff --git a/docs/pipeline.md b/docs/pipeline.md index 8149ca2e..acdd7db8 100644 --- a/docs/pipeline.md +++ b/docs/pipeline.md @@ -208,7 +208,9 @@ The extension has no CUDA or nvCOMP dependency in this configuration. io_uring metadata requirements. macOS uses a POSIX metadata worker pool, with one worker per `metadata_io_concurrency`, and shared POSIX bulk reads. NUMA placement and CPU affinity are unavailable. CUDA defaults off on macOS -and cannot be enabled. +and cannot be enabled. On either platform, `FileReader(workers=...)` and the +legacy `Config.n_io_threads` cannot exceed the online CPU count; larger values +raise `InvalidArgument`. On Linux CUDA defaults on, builds both executors, and requires the CUDA toolkit, nvCOMP, and a runtime NVIDIA driver. GDS requires a CUDA build. diff --git a/python/damacy/__init__.py b/python/damacy/__init__.py index e8c2de09..fa881530 100644 --- a/python/damacy/__init__.py +++ b/python/damacy/__init__.py @@ -501,7 +501,8 @@ class Config: dtype: Destination dtype for assembled batches. lookahead_samples: User-side push-queue depth in samples. Defaults to two full output batches. - n_io_threads: Bulk data IO worker threads (>= 1). Defaults to 64. + n_io_threads: Bulk data IO worker threads, from 1 to the number of + online CPUs. Defaults to 64, so hosts with fewer CPUs must lower it. metadata_io_concurrency: Async metadata request concurrency (>= 1). n_array_meta_cache: LRU cap for zarr-metadata entries. Must be ``>= lookahead_samples + 2 * samples_per_batch`` so the in-flight diff --git a/src/damacy.h b/src/damacy.h index f832fbb0..ce382539 100644 --- a/src/damacy.h +++ b/src/damacy.h @@ -108,9 +108,9 @@ extern "C" // Required: must be in [1, DAMACY_HARD_MAX_SUBSTREAMS_PER_CHUNK]. uint32_t max_substreams_per_chunk; - // Bulk chunk-read worker threads. Wave IO uses this queue. Blocking - // IO workers may exceed the CPU count. Required: must be in - // [1, DAMACY_MAX_IO_THREADS]. + // Bulk chunk-read worker threads. Wave IO uses this queue. Required: + // must be in [1, DAMACY_MAX_IO_THREADS] and no larger than the host's + // online CPU count; creation fails with DAMACY_INVAL otherwise. uint32_t n_io_threads; // Metadata request concurrency for array metadata, shard indexes, and // chunk-layout probes. The Linux metadata path uses this as an io_uring diff --git a/src/pipeline/components.c b/src/pipeline/components.c index 55ece990..e6d5ba94 100644 --- a/src/pipeline/components.c +++ b/src/pipeline/components.c @@ -45,6 +45,15 @@ damacy_file_reader_create(uint32_t workers, *out = NULL; if (!workers || workers > DAMACY_MAX_IO_THREADS || !max_inflight_reads) return DAMACY_INVAL; + int cpus = platform_default_thread_count(); + if (workers > (uint32_t)cpus) { + log_error("file reader workers (n_io_threads)=%u exceeds the %d online " + "CPUs; set it to at most %d", + workers, + cpus, + cpus); + return DAMACY_INVAL; + } struct damacy_reader* reader = calloc(1, sizeof(*reader)); if (!reader) return DAMACY_OOM; diff --git a/tests/test_cpu_pipeline.c b/tests/test_cpu_pipeline.c index 7791b014..5e8e6f3f 100644 --- a/tests/test_cpu_pipeline.c +++ b/tests/test_cpu_pipeline.c @@ -114,6 +114,18 @@ verify_crop(struct damacy_batch* batch, int offset, int unsigned_bits) return 0; } +static int +test_reader_workers_limited_to_cpus(void) +{ + uint32_t cpus = (uint32_t)platform_default_thread_count(); + struct damacy_reader* reader = NULL; + EXPECT(damacy_file_reader_create(cpus + 1, 4, &reader) == DAMACY_INVAL); + EXPECT(!reader); + EXPECT(damacy_file_reader_create(cpus, 4, &reader) == DAMACY_OK); + damacy_reader_destroy(reader); + return 0; +} + static int test_codecs_and_types(void) { @@ -429,5 +441,6 @@ main(void) RUN(test_release_from_other_pipeline); RUN(test_retained_outputs_and_shutdown); RUN(test_bfloat_rounding_and_fill); + RUN(test_reader_workers_limited_to_cpus); return 0; } From d88df0ef91d405396cfc7bea4737bd5200840314 Mon Sep 17 00:00:00 2001 From: Nathan Clack Date: Thu, 24 Sep 2026 22:05:56 +0000 Subject: [PATCH 9/9] format: wrap damacy.h comments --- src/damacy.h | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/src/damacy.h b/src/damacy.h index ce382539..5b677c21 100644 --- a/src/damacy.h +++ b/src/damacy.h @@ -115,7 +115,8 @@ extern "C" // Metadata request concurrency for array metadata, shard indexes, and // chunk-layout probes. The Linux metadata path uses this as an io_uring // request-depth budget; macOS uses this many metadata worker threads. - // Required: must be > 0 and no larger than DAMACY_MAX_METADATA_IO_CONCURRENCY. + // Required: must be > 0 and no larger than + // DAMACY_MAX_METADATA_IO_CONCURRENCY. uint32_t metadata_io_concurrency; uint32_t n_array_meta_cache; @@ -343,8 +344,9 @@ extern "C" } metadata_backend; // Measured metadata-operation latency: submit-to-completion on Linux, // syscall duration on macOS, by kind (stat/open/read/close). Distinct from - // the injected synthetic latency in metadata_latency above. Buckets are log2-scale on ns - // (bucket i: floor(log2(ns)) == i); percentiles are derived from them. + // the injected synthetic latency in metadata_latency above. Buckets are + // log2-scale on ns (bucket i: floor(log2(ns)) == i); percentiles are + // derived from them. struct { uint64_t count;