From 33a6cb6d5b7e16f6cd421a77ea3052313ff98283 Mon Sep 17 00:00:00 2001 From: Nathan Clack Date: Thu, 24 Sep 2026 02:03:01 +0000 Subject: [PATCH 1/5] cpu: buffer 256 chunks per input --- README.md | 2 +- dev/cpu-pipeline.md | 11 +-- docs/pipeline.md | 14 +++- src/CMakeLists.txt | 2 +- src/executor/cpu_executor.c | 89 ++++++++++---------- src/executor/cpu_read_plan.h | 35 ++++++++ tests/test_cpu_executor.c | 157 +++++++++++++++++++++++++++++++++-- tests/test_cpu_pipeline.c | 14 ++-- 8 files changed, 256 insertions(+), 68 deletions(-) create mode 100644 src/executor/cpu_read_plan.h diff --git a/README.md b/README.md index 9b333f05..c9c8f675 100644 --- a/README.md +++ b/README.md @@ -34,7 +34,7 @@ metadata = damacy.ZarrMetadata( planner = damacy.ChunkPlanner(metadata=metadata, limits=damacy.PlanLimits()) executor = damacy.CpuExecutor( reader=chunk_reader, - limits=damacy.CpuLimits(max_memory_bytes=1 << 30, decode_workers=8), + limits=damacy.CpuLimits(max_memory_bytes=3 << 30, decode_workers=8), ) with damacy.Pipeline( diff --git a/dev/cpu-pipeline.md b/dev/cpu-pipeline.md index ce921d8a..a928deef 100644 --- a/dev/cpu-pipeline.md +++ b/dev/cpu-pipeline.md @@ -87,11 +87,12 @@ reuse can be optimized separately. ## Resources and ownership CPU I/O workers and decode workers are configured independently. There are two -input groups and two output buffers. Each input group holds at most -`decode_workers` chunks and `decode_workers * max_encoded_chunk_bytes` encoded -bytes. Merged reads fit these limits, and their count respects the reader -capacity. Decoding uses a bounded per-worker output buffer and zstd context; -memory admission also reserves conservative Blosc workspace. The memory cap +input groups and two output buffers. Each input group holds at most 256 chunks +and `256 * max_encoded_chunk_bytes` encoded bytes, independently of the worker +count. Merged reads contain at most `min(decode_workers, 256)` chunks, and their +count respects the reader capacity. Decoding uses a bounded per-worker output +buffer and zstd context; memory admission also reserves conservative Blosc +workspace. The memory cap covers executor buffers, active read plans, and temporary read-planning scratch. It excludes metadata, reader queues, prepared-plan storage, thread stacks, allocator overhead, and total process RSS. diff --git a/docs/pipeline.md b/docs/pipeline.md index f5b1d804..62649ffb 100644 --- a/docs/pipeline.md +++ b/docs/pipeline.md @@ -29,7 +29,7 @@ planner = damacy.ChunkPlanner( executor = damacy.CpuExecutor( reader=chunk_reader, limits=damacy.CpuLimits( - max_memory_bytes=1 << 30, + max_memory_bytes=3 << 30, decode_workers=8, max_encoded_chunk_bytes=4 << 20, max_decoded_chunk_bytes=2 << 20, @@ -126,8 +126,16 @@ older, stricter cache validation. The CPU executor merges adjacent or overlapping encoded ranges within each shard, then interleaves the reads across shards. It decodes each unique source chunk once per batch, including when several output samples use that chunk. -Both input groups stay within the decode-worker, encoded-byte, and reader -limits; merging does not increase the number of decode workers. +Each of its two encoded-input buffers holds up to 256 chunks, independently +of the decode-worker count. Merged reads contain at most +`min(decode_workers, 256)` chunks; submission also respects the reader limit. +Decoder workspaces remain per worker. + +The two input buffers reserve `512 * max_encoded_chunk_bytes` bytes, plus +per-chunk bookkeeping. The default 4 MiB encoded-chunk bound therefore reserves +2 GiB before decoder workspaces and output buffers. Set the bound to match the +largest encoded chunk expected, and include this reserve in `max_memory_bytes`; +insufficient budgets report `BUDGET`. CPU memory admission includes active read plans, temporary planning scratch, and a conservative allowance for Blosc scratch storage. diff --git a/src/CMakeLists.txt b/src/CMakeLists.txt index 23ee8530..05152a6e 100644 --- a/src/CMakeLists.txt +++ b/src/CMakeLists.txt @@ -278,7 +278,7 @@ if(NOT DAMACY_FUZZ) LINKS pipeline_components plan_builder prefetcher array_meta shard_index metadata_store_async lookahead ) - add_src_lib(cpu_executor SOURCES executor/cpu_executor.c + add_src_lib(cpu_executor SOURCES executor/cpu_executor.c executor/cpu_read_plan.h LINKS pipeline_components prepared_plan dispatch_utils threadpool damacy_stats PkgConfig::ZSTD PkgConfig::BLOSC ) diff --git a/src/executor/cpu_executor.c b/src/executor/cpu_executor.c index 42efb26a..12166de2 100644 --- a/src/executor/cpu_executor.c +++ b/src/executor/cpu_executor.c @@ -7,6 +7,7 @@ #include "damacy_config.h" #include "damacy_stats.h" #include "executor/coalesce.h" +#include "executor/cpu_read_plan.h" #include "threadpool/threadpool.h" #include @@ -20,21 +21,6 @@ enum cpu_slot_state CPU_HELD }; -struct cpu_chunk_read -{ - uint32_t offset; - uint32_t next; -}; - -struct cpu_read_plan -{ - struct read_op* reads; - struct cpu_chunk_read* chunks; - uint32_t* first_chunks; - uint32_t count; - uint64_t bytes; -}; - struct cpu_slot { enum cpu_slot_state state; @@ -93,14 +79,14 @@ struct cpu_executor uint64_t committed; }; -static void +void cpu_read_plan_destroy(struct cpu_read_plan* plan) { free(plan->reads); *plan = (struct cpu_read_plan){ 0 }; } -static enum damacy_status +enum damacy_status cpu_read_plan_build(const struct prepared_plan* plan, const struct damacy_cpu_config* config, uint64_t available, @@ -134,18 +120,21 @@ cpu_read_plan_build(const struct prepared_plan* plan, .file_offset = chunk->offset, .nbytes = chunk->encoded_bytes }; } + uint32_t read_chunks = config->decode_workers; + if (read_chunks > CPU_READ_BUFFER_CHUNKS) + read_chunks = CPU_READ_BUFFER_CHUNKS; uint32_t* read_index = indices + 2 * (size_t)count; uint32_t* offset_in_read = indices + 3 * (size_t)count; reads.count = count; - status = coalesce_reads(reads.reads, - &reads.count, - (uint64_t)config->decode_workers * - config->max_encoded_chunk_bytes, - config->decode_workers, - read_index, - offset_in_read, - indices, - scratch_reads); + status = + coalesce_reads(reads.reads, + &reads.count, + (uint64_t)read_chunks * config->max_encoded_chunk_bytes, + read_chunks, + read_index, + offset_in_read, + indices, + scratch_reads); if (status != DAMACY_OK) goto Done; for (uint32_t i = 0; i < reads.count; ++i) @@ -450,19 +439,27 @@ cpu_start(struct damacy_executor* base, uint64_t workspace = ZSTD_estimateDCtxSize(); uint64_t codec_reserve = workspace + 3ull * self->config.max_decoded_chunk_bytes + (256u << 10); - uint64_t per_worker = + uint64_t per_worker = self->config.max_decoded_chunk_bytes + codec_reserve + + sizeof(struct cpu_worker); + uint64_t per_read_chunk = 2ull * self->config.max_encoded_chunk_bytes + - self->config.max_decoded_chunk_bytes + codec_reserve + - sizeof(struct cpu_worker) + 2 * (sizeof(struct store_read) + sizeof(struct cpu_chunk_input) + sizeof(enum damacy_status) + 2 * sizeof(float) + sizeof(uint64_t)); uint64_t fixed = sizeof(*self) + 2 * sizeof(struct damacy_buffer); - if (per_worker > self->config.max_memory_bytes / workers || - bytes > self->config.max_memory_bytes / 2) - return DAMACY_BUDGET; - uint64_t need = fixed + per_worker * workers + 2 * bytes; - if (need < bytes || need > self->config.max_memory_bytes) + uint64_t budget = self->config.max_memory_bytes; + if (fixed > budget || per_worker > budget / workers || + per_read_chunk > budget / CPU_READ_BUFFER_CHUNKS || bytes > budget / 2 || + CPU_READ_BUFFER_CHUNKS > SIZE_MAX / self->config.max_encoded_chunk_bytes) return DAMACY_BUDGET; + uint64_t need = fixed; + uint64_t parts[] = { per_worker * workers, + per_read_chunk * CPU_READ_BUFFER_CHUNKS, + 2 * bytes }; + for (unsigned i = 0; i < sizeof(parts) / sizeof(*parts); ++i) { + if (parts[i] > budget - need) + return DAMACY_BUDGET; + need += parts[i]; + } self->committed = need; for (unsigned i = 0; i < 2; ++i) { struct damacy_buffer* buffer = calloc(1, sizeof(*buffer)); @@ -477,14 +474,16 @@ cpu_start(struct damacy_executor* base, if (!buffer->data) goto Fail; struct cpu_wave* wave = &self->waves[i]; - wave->input = - malloc((size_t)workers * self->config.max_encoded_chunk_bytes); - wave->reads = calloc((size_t)workers, sizeof(*wave->reads)); - wave->chunks = calloc((size_t)workers, sizeof(*wave->chunks)); - wave->results = calloc((size_t)workers, sizeof(*wave->results)); - wave->decode_ms = calloc((size_t)workers, sizeof(*wave->decode_ms)); - wave->assemble_ms = calloc((size_t)workers, sizeof(*wave->assemble_ms)); - wave->output_bytes = calloc((size_t)workers, sizeof(*wave->output_bytes)); + wave->input = malloc((size_t)CPU_READ_BUFFER_CHUNKS * + self->config.max_encoded_chunk_bytes); + wave->reads = calloc(CPU_READ_BUFFER_CHUNKS, sizeof(*wave->reads)); + wave->chunks = calloc(CPU_READ_BUFFER_CHUNKS, sizeof(*wave->chunks)); + wave->results = calloc(CPU_READ_BUFFER_CHUNKS, sizeof(*wave->results)); + wave->decode_ms = calloc(CPU_READ_BUFFER_CHUNKS, sizeof(*wave->decode_ms)); + wave->assemble_ms = + calloc(CPU_READ_BUFFER_CHUNKS, sizeof(*wave->assemble_ms)); + wave->output_bytes = + calloc(CPU_READ_BUFFER_CHUNKS, sizeof(*wave->output_bytes)); if (!wave->input || !wave->reads || !wave->chunks || !wave->results || !wave->decode_ms || !wave->assemble_ms || !wave->output_bytes) goto Fail; @@ -575,15 +574,15 @@ cpu_dispatch(struct cpu_executor* self, struct cpu_wave* wave, int* changed) uint32_t count = 0; uint32_t n_reads = 0; uint64_t bytes = 0; - uint64_t capacity = (uint64_t)self->config.decode_workers * - self->config.max_encoded_chunk_bytes; + uint64_t capacity = + (uint64_t)CPU_READ_BUFFER_CHUNKS * self->config.max_encoded_chunk_bytes; while (next_read < plan->count) { const struct read_op* read = &plan->reads[next_read]; uint32_t n_chunks = 0; for (uint32_t c = plan->first_chunks[next_read]; c != UINT32_MAX; c = plan->chunks[c].next) ++n_chunks; - if (n_chunks > self->config.decode_workers - count || + if (n_chunks > CPU_READ_BUFFER_CHUNKS - count || read->nbytes > capacity - bytes || (read->nbytes && n_reads == self->reader->max_inflight_reads)) break; diff --git a/src/executor/cpu_read_plan.h b/src/executor/cpu_read_plan.h new file mode 100644 index 00000000..0252d879 --- /dev/null +++ b/src/executor/cpu_read_plan.h @@ -0,0 +1,35 @@ +#pragma once + +#include "damacy_pipeline.h" +#include "executor/dispatch.h" +#include "planner/plan.h" + +enum +{ + CPU_READ_BUFFER_CHUNKS = 256 +}; + +struct cpu_chunk_read +{ + uint32_t offset; + uint32_t next; +}; + +struct cpu_read_plan +{ + struct read_op* reads; + struct cpu_chunk_read* chunks; + uint32_t* first_chunks; + uint32_t count; + uint64_t bytes; +}; + +void +cpu_read_plan_destroy(struct cpu_read_plan* plan); + +enum damacy_status +cpu_read_plan_build(const struct prepared_plan* plan, + const struct damacy_cpu_config* config, + uint64_t available, + uint64_t active_plan_bytes, + struct cpu_read_plan* out); diff --git a/tests/test_cpu_executor.c b/tests/test_cpu_executor.c index c9c18255..c6e44414 100644 --- a/tests/test_cpu_executor.c +++ b/tests/test_cpu_executor.c @@ -1,4 +1,5 @@ #include "damacy_stats.h" +#include "executor/cpu_read_plan.h" #include "expect.h" #include "pipeline/components.h" #include "store/store_internal.h" @@ -8,10 +9,11 @@ enum { - MAX_CHUNKS = 32, - MAX_READS = 64, - SHARD_BYTES = 128, + MAX_CHUNKS = 513, + MAX_READS = 2 * MAX_CHUNKS, + MAX_EVENTS = 32, CHUNK_BYTES = 8, + SHARD_BYTES = MAX_CHUNKS * CHUNK_BYTES, }; struct read_record @@ -33,7 +35,7 @@ struct test_store struct store base; uint16_t data[3][SHARD_BYTES / sizeof(uint16_t)]; struct read_record records[MAX_READS]; - struct pending_reads events[MAX_READS]; + struct pending_reads events[MAX_EVENTS]; unsigned n_reads; unsigned n_events; unsigned pending; @@ -83,7 +85,7 @@ submit_reads(struct store* base, const struct store_read* reads, size_t count) return (struct store_submit_result){ .status = DAMACY_AGAIN }; } if (count > MAX_CHUNKS || count > MAX_READS - store->n_reads || - store->n_events == MAX_READS) + store->n_events == MAX_EVENTS) return (struct store_submit_result){ .status = DAMACY_IO }; struct pending_reads* pending = &store->events[store->n_events++]; pending->count = count; @@ -576,7 +578,7 @@ test_read_errors_and_shutdown(void) else store.held_event = 2; struct damacy_reader reader = { .store = &store.base, - .max_inflight_reads = 2 }; + .max_inflight_reads = 1 }; struct damacy_stats stats = { 0 }; struct damacy_executor* executor = NULL; EXPECT(start_executor(&reader, 2, 4, 8 << 20, &stats, &executor) == 0); @@ -598,6 +600,146 @@ test_read_errors_and_shutdown(void) return 0; } +static int +test_read_buffer_capacity(void) +{ + const uint32_t workers[] = { 1, 2, 32 }; + struct test_chunk chunks[513]; + for (unsigned i = 0; i < 513; ++i) + chunks[i] = (struct test_chunk){ .shard = i % 3, + .offset = (170 - i / 3) * CHUNK_BYTES }; + for (unsigned mode = 0; mode < 3; ++mode) { + struct test_store store; + store_init(&store, 512); + struct damacy_reader reader = { .store = &store.base, + .max_inflight_reads = 512 }; + struct damacy_stats stats = { 0 }; + struct damacy_executor* executor = NULL; + EXPECT(start_executor( + &reader, workers[mode], 513, 32 << 20, &stats, &executor) == 0); + struct prepared_plan* plan = make_plan(chunks, 513); + EXPECT(plan); + EXPECT(executor->ops->submit(executor, plan, 0) == DAMACY_OK); + struct damacy_batch* batch = NULL; + EXPECT(finish_batch(executor, &batch) == 0); + EXPECT(check_output(batch, &store, chunks, 513) == 0); + EXPECT(stats.chunks_dispatched == 513 && stats.waves_emitted == 3); + uint32_t reads_per_shard = (171 + workers[mode] - 1) / workers[mode]; + EXPECT(store.n_reads == 3 * reads_per_shard); + EXPECT(stats.reads_issued == store.n_reads); + for (unsigned i = 0; i < store.n_reads; ++i) { + uint32_t first = (i / 3) * workers[mode]; + uint32_t count = 171 - first; + if (count > workers[mode]) + count = workers[mode]; + EXPECT(store.records[i].shard == i % 3); + EXPECT(store.records[i].offset == (uint64_t)first * CHUNK_BYTES); + EXPECT(store.records[i].bytes == count * CHUNK_BYTES); + } + EXPECT(stats.decode.input_bytes == 513 * CHUNK_BYTES); + EXPECT(stats.io.input_bytes == 513 * CHUNK_BYTES); + damacy_batch_release(batch); + damacy_executor_destroy(executor); + } + return 0; +} + +static int +test_read_plan_many_workers(void) +{ + const uint32_t workers[] = { 256, 257, 1024 }; + struct test_chunk chunks[513]; + for (unsigned i = 0; i < 513; ++i) + chunks[i] = (struct test_chunk){ .offset = (512 - i) * CHUNK_BYTES }; + struct prepared_plan* plan = make_plan(chunks, 513); + EXPECT(plan); + for (unsigned mode = 0; mode < 3; ++mode) { + struct damacy_cpu_config config = { .decode_workers = workers[mode], + .max_encoded_chunk_bytes = CHUNK_BYTES, + .max_decoded_chunk_bytes = CHUNK_BYTES, + .max_memory_bytes = 8 << 20 }; + struct cpu_read_plan reads = { 0 }; + EXPECT(cpu_read_plan_build(plan, &config, 8 << 20, 0, &reads) == DAMACY_OK); + EXPECT(reads.count == 3); + unsigned char seen[513] = { 0 }; + uint32_t total = 0; + for (uint32_t i = 0; i < reads.count; ++i) { + const struct read_op* read = &reads.reads[i]; + EXPECT(!strcmp(read->shard_path, "a")); + EXPECT(read->file_offset == (uint64_t)i * 256 * CHUNK_BYTES); + EXPECT(read->nbytes == (i < 2 ? 256 : 1) * CHUNK_BYTES); + uint32_t count = 0; + for (uint32_t c = reads.first_chunks[i]; c != UINT32_MAX; + c = reads.chunks[c].next) { + EXPECT(c < 513 && !seen[c]); + seen[c] = 1; + EXPECT(read->file_offset + reads.chunks[c].offset == chunks[c].offset); + EXPECT(reads.chunks[c].offset + CHUNK_BYTES <= read->nbytes); + ++count; + } + EXPECT(count == (i < 2 ? 256 : 1)); + total += count; + } + EXPECT(total == 513); + cpu_read_plan_destroy(&reads); + } + prepared_plan_destroy(plan); + return 0; +} + +static int +test_input_memory_budget(void) +{ + struct test_store store; + store_init(&store, 2); + struct damacy_reader reader = { .store = &store.base, + .max_inflight_reads = 2 }; + struct damacy_stats stats = { 0 }; + struct damacy_executor* executor = NULL; + EXPECT(start_executor(&reader, 2, 2, 8 << 20, &stats, &executor) == 0); + executor->ops->stats(executor, &stats); + uint64_t initial = stats.host_bytes_committed; + damacy_executor_destroy(executor); + struct damacy_cpu_config config = { .decode_workers = 2, + .max_encoded_chunk_bytes = + 2 * CHUNK_BYTES, + .max_decoded_chunk_bytes = CHUNK_BYTES, + .max_memory_bytes = 8 << 20 }; + struct damacy_batch_spec output = output_spec(2); + EXPECT(damacy_cpu_executor_create(&reader, &config, &executor) == DAMACY_OK); + EXPECT(executor->ops->start(executor, &output, &stats) == DAMACY_OK); + executor->ops->stats(executor, &stats); + uint64_t required = stats.host_bytes_committed; + EXPECT(required == initial + 512 * CHUNK_BYTES); + damacy_executor_destroy(executor); + const uint64_t budgets[] = { 1, initial, required - 1, required }; + for (unsigned i = 0; i < 4; ++i) { + config.max_memory_bytes = budgets[i]; + EXPECT(damacy_cpu_executor_create(&reader, &config, &executor) == + DAMACY_OK); + EXPECT(executor->ops->start(executor, &output, &stats) == + (i == 3 ? DAMACY_OK : DAMACY_BUDGET)); + executor->ops->stats(executor, &stats); + EXPECT(stats.host_bytes_committed == (i == 3 ? required : 0)); + damacy_executor_destroy(executor); + } + config.max_encoded_chunk_bytes = UINT32_MAX; + config.max_memory_bytes = 8 << 20; + EXPECT(damacy_cpu_executor_create(&reader, &config, &executor) == DAMACY_OK); + EXPECT(executor->ops->start(executor, &output, &stats) == DAMACY_BUDGET); + damacy_executor_destroy(executor); + if (SIZE_MAX > UINT32_MAX) { + config.max_encoded_chunk_bytes = CHUNK_BYTES; + config.max_memory_bytes = UINT64_MAX; + output.sample_shape[0] = INT64_MAX / 8; + EXPECT(damacy_cpu_executor_create(&reader, &config, &executor) == + DAMACY_OK); + EXPECT(executor->ops->start(executor, &output, &stats) == DAMACY_BUDGET); + damacy_executor_destroy(executor); + } + return 0; +} + int main(void) { @@ -610,5 +752,8 @@ main(void) RUN(test_plan_memory_budget); RUN(test_plan_memory_retry); RUN(test_read_errors_and_shutdown); + RUN(test_read_buffer_capacity); + RUN(test_read_plan_many_workers); + RUN(test_input_memory_budget); return 0; } diff --git a/tests/test_cpu_pipeline.c b/tests/test_cpu_pipeline.c index 5e8e6f3f..92b0150e 100644 --- a/tests/test_cpu_pipeline.c +++ b/tests/test_cpu_pipeline.c @@ -52,13 +52,13 @@ create_components(struct components* c) .max_shards_per_sample = 4, .max_plan_bytes = 1 << 20 }, &c->planner) == DAMACY_OK); - EXPECT(damacy_cpu_executor_create(c->reader, - &(struct damacy_cpu_config){ - .decode_workers = 2, - .max_encoded_chunk_bytes = (1 << 20) + 1, - .max_decoded_chunk_bytes = 1 << 20, - .max_memory_bytes = 32 << 20 }, - &c->executor) == DAMACY_OK); + EXPECT(damacy_cpu_executor_create( + c->reader, + &(struct damacy_cpu_config){ .decode_workers = 2, + .max_encoded_chunk_bytes = 1024, + .max_decoded_chunk_bytes = 1 << 20, + .max_memory_bytes = 32 << 20 }, + &c->executor) == DAMACY_OK); return 0; } From 43f91607d803eb57edeecc182b78a2d0f36cf67f Mon Sep 17 00:00:00 2001 From: Nathan Clack Date: Thu, 24 Sep 2026 21:26:24 +0000 Subject: [PATCH 2/5] cpu: set chunks per input buffer --- bench/main.c | 12 ++- dev/cpu-pipeline.md | 10 +- docs/pipeline.md | 19 ++-- python/damacy/__init__.py | 8 ++ python/damacy/_components.c | 5 +- python/damacy/_native.pyi | 8 +- python/tests/test_components.py | 3 + src/damacy_pipeline.h | 1 + src/executor/cpu_executor.c | 56 ++++++----- src/executor/cpu_read_plan.h | 5 - tests/test_cpu_executor.c | 165 +++++++++++++++++++++++++------- tests/test_cpu_pipeline.c | 3 +- 12 files changed, 211 insertions(+), 84 deletions(-) diff --git a/bench/main.c b/bench/main.c index 0d3461d9..afdcaeee 100644 --- a/bench/main.c +++ b/bench/main.c @@ -120,6 +120,7 @@ struct scenario uint64_t max_gpu_memory_bytes; uint64_t max_cpu_memory_bytes; uint32_t decode_workers; + uint32_t chunks_per_input_buffer; int cpu; uint32_t max_chunk_uncompressed_bytes; // 0 → tuning_defaults() baseline uint64_t max_read_op_bytes; // 0 → tuning_defaults() baseline @@ -451,6 +452,10 @@ parse_scenario(struct cslice src, struct scenario* sc) static const struct json_query p_workers[] = { { QUERY_KEY, .key = "pipeline" }, { QUERY_KEY, .key = "decode_workers" } }; + static const struct json_query p_buffer_chunks[] = { + { QUERY_KEY, .key = "pipeline" }, + { QUERY_KEY, .key = "chunks_per_input_buffer" } + }; read_uint_opt(src, p_cpu, countof(p_cpu), &v, 0); if (v > (UINT64_MAX >> 20) || (sc->cpu && !v)) return 1; @@ -459,6 +464,10 @@ parse_scenario(struct cslice src, struct scenario* sc) if (!v || v > UINT32_MAX) return 1; sc->decode_workers = (uint32_t)v; + read_uint_opt(src, p_buffer_chunks, countof(p_buffer_chunks), &v, 256); + if (!v || v > UINT32_MAX) + return 1; + sc->chunks_per_input_buffer = (uint32_t)v; read_uint_opt(src, p_g, countof(p_g), &v, 0); sc->max_gpu_memory_bytes = v << 20; read_uint_opt(src, p_c, countof(p_c), &v, 0); @@ -1134,7 +1143,8 @@ pipeline_create(const struct scenario* scenario, .decode_workers = scenario->decode_workers, .max_encoded_chunk_bytes = (uint32_t)cfg->tuning.max_read_op_bytes, .max_decoded_chunk_bytes = cfg->tuning.max_chunk_uncompressed_bytes, - .max_memory_bytes = scenario->max_cpu_memory_bytes }, + .max_memory_bytes = scenario->max_cpu_memory_bytes, + .chunks_per_input_buffer = scenario->chunks_per_input_buffer }, &pipeline->executor); if (status != DAMACY_OK) return status; diff --git a/dev/cpu-pipeline.md b/dev/cpu-pipeline.md index a928deef..c7c52cbf 100644 --- a/dev/cpu-pipeline.md +++ b/dev/cpu-pipeline.md @@ -87,10 +87,12 @@ reuse can be optimized separately. ## Resources and ownership CPU I/O workers and decode workers are configured independently. There are two -input groups and two output buffers. Each input group holds at most 256 chunks -and `256 * max_encoded_chunk_bytes` encoded bytes, independently of the worker -count. Merged reads contain at most `min(decode_workers, 256)` chunks, and their -count respects the reader capacity. Decoding uses a bounded per-worker output +input groups and two output buffers. Each input group holds at most +`chunks_per_input_buffer` chunks and +`chunks_per_input_buffer * max_encoded_chunk_bytes` encoded bytes. The setting +is at least `decode_workers`, so every worker can get a chunk. Merged reads +contain at most `chunks_per_input_buffer` chunks, and their count respects the +reader capacity. Decoding uses a bounded per-worker output buffer and zstd context; memory admission also reserves conservative Blosc workspace. The memory cap covers executor buffers, active read plans, and temporary read-planning diff --git a/docs/pipeline.md b/docs/pipeline.md index 62649ffb..72bfce8e 100644 --- a/docs/pipeline.md +++ b/docs/pipeline.md @@ -33,6 +33,7 @@ executor = damacy.CpuExecutor( decode_workers=8, max_encoded_chunk_bytes=4 << 20, max_decoded_chunk_bytes=2 << 20, + chunks_per_input_buffer=256, ), ) output = damacy.BatchSpec(samples=2, shape=(64, 256, 256), dtype="f32") @@ -112,6 +113,7 @@ capacities. Defaults come from the Python value objects. | `PlanLimits.max_plan_bytes` | Owned storage per prepared plan. | | `CpuLimits.max_memory_bytes` | CPU executor's buffers, codec workspace, active read plans, and temporary read-planning scratch. | | `CpuLimits.decode_workers` | Total decoding/assembly workers, including the calling scheduler thread. | +| `CpuLimits.chunks_per_input_buffer` | Chunks each of the two encoded-input buffers holds; from `decode_workers` to 16384. | | `FileReader.workers` | Bulk I/O workers, separate from decoding workers. | | `FileReader.max_inflight_reads` | Bulk read capacity; execution respects this bound and retries saturation. | | `CudaLimits` | GPU memory and execution geometry, plus the CUDA codec-layout cache capacity. | @@ -126,16 +128,17 @@ older, stricter cache validation. The CPU executor merges adjacent or overlapping encoded ranges within each shard, then interleaves the reads across shards. It decodes each unique source chunk once per batch, including when several output samples use that chunk. -Each of its two encoded-input buffers holds up to 256 chunks, independently -of the decode-worker count. Merged reads contain at most -`min(decode_workers, 256)` chunks; submission also respects the reader limit. +Each of its two encoded-input buffers holds `chunks_per_input_buffer` chunks +(default 256). A merged read contains at most that many chunks, so reads merge +even with one decode worker; submission also respects the reader limit. Decoder workspaces remain per worker. -The two input buffers reserve `512 * max_encoded_chunk_bytes` bytes, plus -per-chunk bookkeeping. The default 4 MiB encoded-chunk bound therefore reserves -2 GiB before decoder workspaces and output buffers. Set the bound to match the -largest encoded chunk expected, and include this reserve in `max_memory_bytes`; -insufficient budgets report `BUDGET`. +The two input buffers reserve +`2 * chunks_per_input_buffer * max_encoded_chunk_bytes` bytes, plus per-chunk +bookkeeping. With the defaults, 256 chunks and a 4 MiB encoded-chunk bound, +that is 2 GiB before decoder workspaces and output buffers. Set the bound to +match the largest encoded chunk expected, and include this reserve in +`max_memory_bytes`; insufficient budgets report `BUDGET`. CPU memory admission includes active read plans, temporary planning scratch, and a conservative allowance for Blosc scratch storage. diff --git a/python/damacy/__init__.py b/python/damacy/__init__.py index fa881530..8aec0c46 100644 --- a/python/damacy/__init__.py +++ b/python/damacy/__init__.py @@ -826,12 +826,19 @@ class CpuLimits: decode_workers: int = 8 max_encoded_chunk_bytes: int = 4 << 20 max_decoded_chunk_bytes: int = 2 << 20 + chunks_per_input_buffer: int = 256 def __post_init__(self) -> None: _positive_int(self.max_memory_bytes, "max_memory_bytes", (1 << 64) - 1) _positive_int(self.decode_workers, "decode_workers", _native.MAX_IO_THREADS) _positive_int(self.max_encoded_chunk_bytes, "max_encoded_chunk_bytes") _positive_int(self.max_decoded_chunk_bytes, "max_decoded_chunk_bytes") + _positive_int(self.chunks_per_input_buffer, "chunks_per_input_buffer", 16384) + if self.chunks_per_input_buffer < self.decode_workers: + raise ValueError( + f"chunks_per_input_buffer ({self.chunks_per_input_buffer}) must be " + f"at least decode_workers ({self.decode_workers})" + ) @dataclass(frozen=True, slots=True) @@ -962,6 +969,7 @@ def __init__(self, *, reader: FileReader, limits: CpuLimits) -> None: limits.max_encoded_chunk_bytes, limits.max_decoded_chunk_bytes, limits.max_memory_bytes, + limits.chunks_per_input_buffer, ) diff --git a/python/damacy/_components.c b/python/damacy/_components.c index 16677b4d..d357a30b 100644 --- a/python/damacy/_components.c +++ b/python/damacy/_components.c @@ -186,12 +186,13 @@ create_cpu_executor(PyObject* self, PyObject* args) unsigned long long bytes; struct damacy_cpu_config config; if (!PyArg_ParseTuple(args, - "OIIIK", + "OIIIKI", &reader_object, &config.decode_workers, &config.max_encoded_chunk_bytes, &config.max_decoded_chunk_bytes, - &bytes)) + &bytes, + &config.chunks_per_input_buffer)) return NULL; config.max_memory_bytes = bytes; struct damacy_reader* reader = component_value(reader_object, READER); diff --git a/python/damacy/_native.pyi b/python/damacy/_native.pyi index 423487ef..b7a5ad24 100644 --- a/python/damacy/_native.pyi +++ b/python/damacy/_native.pyi @@ -214,7 +214,13 @@ def create_planner( /, ) -> object: ... def create_cpu_executor( - reader: object, workers: int, max_encoded: int, max_decoded: int, max_memory: int, / + reader: object, + workers: int, + max_encoded: int, + max_decoded: int, + max_memory: int, + chunks_per_input_buffer: int, + /, ) -> object: ... def create_cuda_executor( reader: object, diff --git a/python/tests/test_components.py b/python/tests/test_components.py index 262f9bb8..7c9c02d5 100644 --- a/python/tests/test_components.py +++ b/python/tests/test_components.py @@ -284,6 +284,9 @@ def test_invalid_limits(): lambda: damacy.CpuLimits(0), lambda: damacy.PlanLimits(max_chunks=0), lambda: damacy.CpuLimits(1 << 20, decode_workers=-1), + lambda: damacy.CpuLimits(1 << 20, chunks_per_input_buffer=0), + lambda: damacy.CpuLimits(1 << 20, decode_workers=8, chunks_per_input_buffer=7), + lambda: damacy.CpuLimits(1 << 20, chunks_per_input_buffer=16385), ]: with pytest.raises(ValueError): make() diff --git a/src/damacy_pipeline.h b/src/damacy_pipeline.h index bab1b71f..7f09c6cf 100644 --- a/src/damacy_pipeline.h +++ b/src/damacy_pipeline.h @@ -47,6 +47,7 @@ extern "C" uint32_t max_encoded_chunk_bytes; uint32_t max_decoded_chunk_bytes; uint64_t max_memory_bytes; + uint32_t chunks_per_input_buffer; }; struct damacy_cuda_config diff --git a/src/executor/cpu_executor.c b/src/executor/cpu_executor.c index 12166de2..d964c3d5 100644 --- a/src/executor/cpu_executor.c +++ b/src/executor/cpu_executor.c @@ -8,6 +8,7 @@ #include "damacy_stats.h" #include "executor/coalesce.h" #include "executor/cpu_read_plan.h" +#include "log/log.h" #include "threadpool/threadpool.h" #include @@ -120,9 +121,7 @@ cpu_read_plan_build(const struct prepared_plan* plan, .file_offset = chunk->offset, .nbytes = chunk->encoded_bytes }; } - uint32_t read_chunks = config->decode_workers; - if (read_chunks > CPU_READ_BUFFER_CHUNKS) - read_chunks = CPU_READ_BUFFER_CHUNKS; + uint32_t read_chunks = config->chunks_per_input_buffer; uint32_t* read_index = indices + 2 * (size_t)count; uint32_t* offset_in_read = indices + 3 * (size_t)count; reads.count = count; @@ -446,21 +445,15 @@ cpu_start(struct damacy_executor* base, 2 * (sizeof(struct store_read) + sizeof(struct cpu_chunk_input) + sizeof(enum damacy_status) + 2 * sizeof(float) + sizeof(uint64_t)); uint64_t fixed = sizeof(*self) + 2 * sizeof(struct damacy_buffer); + uint64_t chunks = self->config.chunks_per_input_buffer; + uint64_t input = per_read_chunk * chunks; + uint64_t decode = per_worker * workers; uint64_t budget = self->config.max_memory_bytes; - if (fixed > budget || per_worker > budget / workers || - per_read_chunk > budget / CPU_READ_BUFFER_CHUNKS || bytes > budget / 2 || - CPU_READ_BUFFER_CHUNKS > SIZE_MAX / self->config.max_encoded_chunk_bytes) + uint64_t need = fixed + input + decode; + if (chunks > SIZE_MAX / self->config.max_encoded_chunk_bytes || + need > budget || bytes > (budget - need) / 2) return DAMACY_BUDGET; - uint64_t need = fixed; - uint64_t parts[] = { per_worker * workers, - per_read_chunk * CPU_READ_BUFFER_CHUNKS, - 2 * bytes }; - for (unsigned i = 0; i < sizeof(parts) / sizeof(*parts); ++i) { - if (parts[i] > budget - need) - return DAMACY_BUDGET; - need += parts[i]; - } - self->committed = need; + self->committed = need + 2 * bytes; for (unsigned i = 0; i < 2; ++i) { struct damacy_buffer* buffer = calloc(1, sizeof(*buffer)); if (!buffer) @@ -474,16 +467,13 @@ cpu_start(struct damacy_executor* base, if (!buffer->data) goto Fail; struct cpu_wave* wave = &self->waves[i]; - wave->input = malloc((size_t)CPU_READ_BUFFER_CHUNKS * - self->config.max_encoded_chunk_bytes); - wave->reads = calloc(CPU_READ_BUFFER_CHUNKS, sizeof(*wave->reads)); - wave->chunks = calloc(CPU_READ_BUFFER_CHUNKS, sizeof(*wave->chunks)); - wave->results = calloc(CPU_READ_BUFFER_CHUNKS, sizeof(*wave->results)); - wave->decode_ms = calloc(CPU_READ_BUFFER_CHUNKS, sizeof(*wave->decode_ms)); - wave->assemble_ms = - calloc(CPU_READ_BUFFER_CHUNKS, sizeof(*wave->assemble_ms)); - wave->output_bytes = - calloc(CPU_READ_BUFFER_CHUNKS, sizeof(*wave->output_bytes)); + wave->input = malloc((size_t)chunks * self->config.max_encoded_chunk_bytes); + wave->reads = calloc((size_t)chunks, sizeof(*wave->reads)); + wave->chunks = calloc((size_t)chunks, sizeof(*wave->chunks)); + wave->results = calloc((size_t)chunks, sizeof(*wave->results)); + wave->decode_ms = calloc((size_t)chunks, sizeof(*wave->decode_ms)); + wave->assemble_ms = calloc((size_t)chunks, sizeof(*wave->assemble_ms)); + wave->output_bytes = calloc((size_t)chunks, sizeof(*wave->output_bytes)); if (!wave->input || !wave->reads || !wave->chunks || !wave->results || !wave->decode_ms || !wave->assemble_ms || !wave->output_bytes) goto Fail; @@ -574,16 +564,16 @@ cpu_dispatch(struct cpu_executor* self, struct cpu_wave* wave, int* changed) uint32_t count = 0; uint32_t n_reads = 0; uint64_t bytes = 0; + uint32_t chunk_capacity = self->config.chunks_per_input_buffer; uint64_t capacity = - (uint64_t)CPU_READ_BUFFER_CHUNKS * self->config.max_encoded_chunk_bytes; + (uint64_t)chunk_capacity * self->config.max_encoded_chunk_bytes; while (next_read < plan->count) { const struct read_op* read = &plan->reads[next_read]; uint32_t n_chunks = 0; for (uint32_t c = plan->first_chunks[next_read]; c != UINT32_MAX; c = plan->chunks[c].next) ++n_chunks; - if (n_chunks > CPU_READ_BUFFER_CHUNKS - count || - read->nbytes > capacity - bytes || + if (n_chunks > chunk_capacity - count || read->nbytes > capacity - bytes || (read->nbytes && n_reads == self->reader->max_inflight_reads)) break; for (uint32_t c = plan->first_chunks[next_read]; c != UINT32_MAX; @@ -755,6 +745,14 @@ damacy_cpu_executor_create(struct damacy_reader* reader, !config->max_encoded_chunk_bytes || !config->max_decoded_chunk_bytes || !config->max_memory_bytes) return DAMACY_INVAL; + if (config->chunks_per_input_buffer < config->decode_workers || + config->chunks_per_input_buffer > DAMACY_MAX_CHUNKS_PER_BATCH) { + log_error("chunks_per_input_buffer=%u out of range (decode_workers=%u..%u)", + config->chunks_per_input_buffer, + config->decode_workers, + DAMACY_MAX_CHUNKS_PER_BATCH); + return DAMACY_INVAL; + } struct cpu_executor* self = calloc(1, sizeof(*self)); if (!self) return DAMACY_OOM; diff --git a/src/executor/cpu_read_plan.h b/src/executor/cpu_read_plan.h index 0252d879..983a9ef9 100644 --- a/src/executor/cpu_read_plan.h +++ b/src/executor/cpu_read_plan.h @@ -4,11 +4,6 @@ #include "executor/dispatch.h" #include "planner/plan.h" -enum -{ - CPU_READ_BUFFER_CHUNKS = 256 -}; - struct cpu_chunk_read { uint32_t offset; diff --git a/tests/test_cpu_executor.c b/tests/test_cpu_executor.c index c6e44414..0fb28a47 100644 --- a/tests/test_cpu_executor.c +++ b/tests/test_cpu_executor.c @@ -216,6 +216,7 @@ make_plan(const struct test_chunk* chunks, uint32_t count) static int start_executor(struct damacy_reader* reader, uint32_t workers, + uint32_t buffer_chunks, uint32_t count, uint64_t budget, struct damacy_stats* stats, @@ -224,7 +225,9 @@ start_executor(struct damacy_reader* reader, struct damacy_cpu_config config = { .decode_workers = workers, .max_encoded_chunk_bytes = CHUNK_BYTES, .max_decoded_chunk_bytes = CHUNK_BYTES, - .max_memory_bytes = budget }; + .max_memory_bytes = budget, + .chunks_per_input_buffer = + buffer_chunks }; EXPECT(damacy_cpu_executor_create(reader, &config, out) == DAMACY_OK); struct damacy_batch_spec output = output_spec(count); EXPECT((*out)->ops->start(*out, &output, stats) == DAMACY_OK); @@ -283,7 +286,7 @@ test_merge_and_shard_order(void) .max_inflight_reads = 4 }; struct damacy_stats stats = { 0 }; struct damacy_executor* executor = NULL; - EXPECT(start_executor(&reader, 4, 12, 8 << 20, &stats, &executor) == 0); + EXPECT(start_executor(&reader, 4, 4, 12, 8 << 20, &stats, &executor) == 0); struct prepared_plan* plan = make_plan(chunks, 12); EXPECT(plan); EXPECT(plan->chunks[0].path != plan->chunks[1].path); @@ -319,7 +322,7 @@ test_reader_limit_and_retry(void) .max_inflight_reads = 1 }; struct damacy_stats stats = { 0 }; struct damacy_executor* executor = NULL; - EXPECT(start_executor(&reader, 2, 4, 8 << 20, &stats, &executor) == 0); + EXPECT(start_executor(&reader, 2, 2, 4, 8 << 20, &stats, &executor) == 0); struct prepared_plan* plan = make_plan(chunks, 4); EXPECT(plan); EXPECT(executor->ops->submit(executor, plan, 0) == DAMACY_OK); @@ -345,7 +348,7 @@ test_overlapping_ranges(void) .max_inflight_reads = 2 }; struct damacy_stats stats = { 0 }; struct damacy_executor* executor = NULL; - EXPECT(start_executor(&reader, 2, 2, 8 << 20, &stats, &executor) == 0); + EXPECT(start_executor(&reader, 2, 2, 2, 8 << 20, &stats, &executor) == 0); struct prepared_plan* plan = make_plan(chunks, 2); EXPECT(plan); EXPECT(executor->ops->submit(executor, plan, 0) == DAMACY_OK); @@ -376,7 +379,7 @@ test_single_worker_and_fills(void) .max_inflight_reads = 1 }; struct damacy_stats stats = { 0 }; struct damacy_executor* executor = NULL; - EXPECT(start_executor(&reader, 1, 4, 8 << 20, &stats, &executor) == 0); + EXPECT(start_executor(&reader, 1, 1, 4, 8 << 20, &stats, &executor) == 0); struct prepared_plan* plan = make_plan(chunks, 4); EXPECT(plan); EXPECT(executor->ops->submit(executor, plan, 0) == DAMACY_OK); @@ -405,7 +408,7 @@ test_fills_share_merged_input(void) .max_inflight_reads = 2 }; struct damacy_stats stats = { 0 }; struct damacy_executor* executor = NULL; - EXPECT(start_executor(&reader, 6, 6, 8 << 20, &stats, &executor) == 0); + EXPECT(start_executor(&reader, 6, 6, 6, 8 << 20, &stats, &executor) == 0); struct prepared_plan* plan = make_plan(chunks, 6); EXPECT(plan); EXPECT(executor->ops->submit(executor, plan, 0) == DAMACY_OK); @@ -437,7 +440,7 @@ test_batch_order_and_retained_output(void) .max_inflight_reads = 2 }; struct damacy_stats stats = { 0 }; struct damacy_executor* executor = NULL; - EXPECT(start_executor(&reader, 2, 2, 8 << 20, &stats, &executor) == 0); + EXPECT(start_executor(&reader, 2, 2, 2, 8 << 20, &stats, &executor) == 0); for (unsigned i = 0; i < 2; ++i) { struct prepared_plan* plan = make_plan(chunks, 2); EXPECT(plan); @@ -484,7 +487,7 @@ test_plan_memory_budget(void) .max_inflight_reads = 2 }; struct damacy_stats stats = { 0 }; struct damacy_executor* executor = NULL; - EXPECT(start_executor(&reader, 2, 2, 8 << 20, &stats, &executor) == 0); + EXPECT(start_executor(&reader, 2, 2, 2, 8 << 20, &stats, &executor) == 0); executor->ops->stats(executor, &stats); uint64_t initial = stats.host_bytes_committed; struct prepared_plan* plan = make_plan(chunks, 2); @@ -500,7 +503,7 @@ test_plan_memory_budget(void) damacy_batch_release(batch); damacy_executor_destroy(executor); executor = NULL; - EXPECT(start_executor(&reader, 2, 2, with_plan, &stats, &executor) == 0); + EXPECT(start_executor(&reader, 2, 2, 2, with_plan, &stats, &executor) == 0); plan = make_plan(chunks, 2); EXPECT(plan); EXPECT(executor->ops->submit(executor, plan, 0) == DAMACY_BUDGET); @@ -521,7 +524,7 @@ test_plan_memory_retry(void) .max_inflight_reads = 2 }; struct damacy_stats stats = { 0 }; struct damacy_executor* executor = NULL; - EXPECT(start_executor(&reader, 2, 2, 8 << 20, &stats, &executor) == 0); + EXPECT(start_executor(&reader, 2, 2, 2, 8 << 20, &stats, &executor) == 0); executor->ops->stats(executor, &stats); uint64_t initial = stats.host_bytes_committed; struct prepared_plan* plan = make_plan(chunks, 2); @@ -534,7 +537,8 @@ test_plan_memory_retry(void) for (unsigned i = 1; i <= 16; ++i) { executor = NULL; EXPECT(start_executor( - &reader, 2, 2, initial + i * plan_bytes, &stats, &executor) == 0); + &reader, 2, 2, 2, initial + i * plan_bytes, &stats, &executor) == + 0); plan = make_plan(chunks, 2); EXPECT(plan); enum damacy_status status = executor->ops->submit(executor, plan, 0); @@ -581,7 +585,7 @@ test_read_errors_and_shutdown(void) .max_inflight_reads = 1 }; struct damacy_stats stats = { 0 }; struct damacy_executor* executor = NULL; - EXPECT(start_executor(&reader, 2, 4, 8 << 20, &stats, &executor) == 0); + EXPECT(start_executor(&reader, 2, 2, 4, 8 << 20, &stats, &executor) == 0); struct prepared_plan* plan = make_plan(chunks, 4); EXPECT(plan); EXPECT(executor->ops->submit(executor, plan, 0) == DAMACY_OK); @@ -604,10 +608,12 @@ static int test_read_buffer_capacity(void) { const uint32_t workers[] = { 1, 2, 32 }; + const uint32_t buffer_chunks = 256; struct test_chunk chunks[513]; for (unsigned i = 0; i < 513; ++i) - chunks[i] = (struct test_chunk){ .shard = i % 3, - .offset = (170 - i / 3) * CHUNK_BYTES }; + chunks[i] = + (struct test_chunk){ .shard = i % 3, + .offset = 2 * (170 - i / 3) * CHUNK_BYTES }; for (unsigned mode = 0; mode < 3; ++mode) { struct test_store store; store_init(&store, 512); @@ -615,8 +621,13 @@ test_read_buffer_capacity(void) .max_inflight_reads = 512 }; struct damacy_stats stats = { 0 }; struct damacy_executor* executor = NULL; - EXPECT(start_executor( - &reader, workers[mode], 513, 32 << 20, &stats, &executor) == 0); + EXPECT(start_executor(&reader, + workers[mode], + buffer_chunks, + 513, + 32 << 20, + &stats, + &executor) == 0); struct prepared_plan* plan = make_plan(chunks, 513); EXPECT(plan); EXPECT(executor->ops->submit(executor, plan, 0) == DAMACY_OK); @@ -624,19 +635,64 @@ test_read_buffer_capacity(void) EXPECT(finish_batch(executor, &batch) == 0); EXPECT(check_output(batch, &store, chunks, 513) == 0); EXPECT(stats.chunks_dispatched == 513 && stats.waves_emitted == 3); - uint32_t reads_per_shard = (171 + workers[mode] - 1) / workers[mode]; + EXPECT(store.n_events == 3); + EXPECT(store.events[0].count == buffer_chunks); + EXPECT(store.events[1].count == buffer_chunks); + EXPECT(store.events[2].count == 513 - 2 * buffer_chunks); + EXPECT(store.n_reads == 513 && stats.reads_issued == 513); + for (unsigned i = 0; i < store.n_reads; ++i) { + EXPECT(store.records[i].shard == i % 3); + EXPECT(store.records[i].offset == 2 * (uint64_t)(i / 3) * CHUNK_BYTES); + EXPECT(store.records[i].bytes == CHUNK_BYTES); + } + EXPECT(stats.decode.input_bytes == 513 * CHUNK_BYTES); + EXPECT(stats.io.input_bytes == 513 * CHUNK_BYTES); + damacy_batch_release(batch); + damacy_executor_destroy(executor); + } + return 0; +} + +static int +test_merged_read_limit(void) +{ + const uint32_t buffer_chunks[] = { 64, 256 }; + struct test_chunk chunks[513]; + for (unsigned i = 0; i < 513; ++i) + chunks[i] = (struct test_chunk){ .shard = i % 3, + .offset = (170 - i / 3) * CHUNK_BYTES }; + for (unsigned mode = 0; mode < 2; ++mode) { + uint32_t limit = buffer_chunks[mode]; + struct test_store store; + store_init(&store, 512); + struct damacy_reader reader = { .store = &store.base, + .max_inflight_reads = 512 }; + struct damacy_stats stats = { 0 }; + struct damacy_executor* executor = NULL; + EXPECT( + start_executor(&reader, 1, limit, 513, 32 << 20, &stats, &executor) == 0); + struct prepared_plan* plan = make_plan(chunks, 513); + EXPECT(plan); + EXPECT(executor->ops->submit(executor, plan, 0) == DAMACY_OK); + struct damacy_batch* batch = NULL; + EXPECT(finish_batch(executor, &batch) == 0); + EXPECT(check_output(batch, &store, chunks, 513) == 0); + uint32_t reads_per_shard = (171 + limit - 1) / limit; EXPECT(store.n_reads == 3 * reads_per_shard); + EXPECT(store.n_events == store.n_reads); EXPECT(stats.reads_issued == store.n_reads); + EXPECT(stats.waves_emitted == store.n_reads); for (unsigned i = 0; i < store.n_reads; ++i) { - uint32_t first = (i / 3) * workers[mode]; + uint32_t first = (i / 3) * limit; uint32_t count = 171 - first; - if (count > workers[mode]) - count = workers[mode]; + if (count > limit) + count = limit; + EXPECT(store.events[i].count == 1); EXPECT(store.records[i].shard == i % 3); EXPECT(store.records[i].offset == (uint64_t)first * CHUNK_BYTES); EXPECT(store.records[i].bytes == count * CHUNK_BYTES); } - EXPECT(stats.decode.input_bytes == 513 * CHUNK_BYTES); + EXPECT(stats.chunks_dispatched == 513); EXPECT(stats.io.input_bytes == 513 * CHUNK_BYTES); damacy_batch_release(batch); damacy_executor_destroy(executor); @@ -645,29 +701,34 @@ test_read_buffer_capacity(void) } static int -test_read_plan_many_workers(void) +test_read_plan_large_buffers(void) { - const uint32_t workers[] = { 256, 257, 1024 }; + const uint32_t buffer_chunks[] = { 256, 257, 1024 }; struct test_chunk chunks[513]; for (unsigned i = 0; i < 513; ++i) chunks[i] = (struct test_chunk){ .offset = (512 - i) * CHUNK_BYTES }; struct prepared_plan* plan = make_plan(chunks, 513); EXPECT(plan); for (unsigned mode = 0; mode < 3; ++mode) { - struct damacy_cpu_config config = { .decode_workers = workers[mode], + uint32_t limit = buffer_chunks[mode]; + struct damacy_cpu_config config = { .decode_workers = 1, .max_encoded_chunk_bytes = CHUNK_BYTES, .max_decoded_chunk_bytes = CHUNK_BYTES, - .max_memory_bytes = 8 << 20 }; + .max_memory_bytes = 8 << 20, + .chunks_per_input_buffer = limit }; struct cpu_read_plan reads = { 0 }; EXPECT(cpu_read_plan_build(plan, &config, 8 << 20, 0, &reads) == DAMACY_OK); - EXPECT(reads.count == 3); + EXPECT(reads.count == (513 + limit - 1) / limit); unsigned char seen[513] = { 0 }; uint32_t total = 0; for (uint32_t i = 0; i < reads.count; ++i) { const struct read_op* read = &reads.reads[i]; + uint32_t expected = 513 - i * limit; + if (expected > limit) + expected = limit; EXPECT(!strcmp(read->shard_path, "a")); - EXPECT(read->file_offset == (uint64_t)i * 256 * CHUNK_BYTES); - EXPECT(read->nbytes == (i < 2 ? 256 : 1) * CHUNK_BYTES); + EXPECT(read->file_offset == (uint64_t)i * limit * CHUNK_BYTES); + EXPECT(read->nbytes == expected * CHUNK_BYTES); uint32_t count = 0; for (uint32_t c = reads.first_chunks[i]; c != UINT32_MAX; c = reads.chunks[c].next) { @@ -677,7 +738,7 @@ test_read_plan_many_workers(void) EXPECT(reads.chunks[c].offset + CHUNK_BYTES <= read->nbytes); ++count; } - EXPECT(count == (i < 2 ? 256 : 1)); + EXPECT(count == expected); total += count; } EXPECT(total == 513); @@ -687,6 +748,41 @@ test_read_plan_many_workers(void) return 0; } +static int +test_input_buffer_validation(void) +{ + struct test_store store; + store_init(&store, 2); + struct damacy_reader reader = { .store = &store.base, + .max_inflight_reads = 2 }; + const struct + { + uint32_t workers; + uint32_t buffer_chunks; + enum damacy_status status; + } cases[] = { + { 8, 0, DAMACY_INVAL }, + { 8, 7, DAMACY_INVAL }, + { 8, 8, DAMACY_OK }, + { 1, DAMACY_MAX_CHUNKS_PER_BATCH, DAMACY_OK }, + { 1, DAMACY_MAX_CHUNKS_PER_BATCH + 1, DAMACY_INVAL }, + }; + for (unsigned i = 0; i < sizeof(cases) / sizeof(*cases); ++i) { + struct damacy_cpu_config config = { .decode_workers = cases[i].workers, + .max_encoded_chunk_bytes = CHUNK_BYTES, + .max_decoded_chunk_bytes = CHUNK_BYTES, + .max_memory_bytes = 8 << 20, + .chunks_per_input_buffer = + cases[i].buffer_chunks }; + struct damacy_executor* executor = NULL; + EXPECT(damacy_cpu_executor_create(&reader, &config, &executor) == + cases[i].status); + EXPECT((executor != NULL) == (cases[i].status == DAMACY_OK)); + damacy_executor_destroy(executor); + } + return 0; +} + static int test_input_memory_budget(void) { @@ -696,7 +792,7 @@ test_input_memory_budget(void) .max_inflight_reads = 2 }; struct damacy_stats stats = { 0 }; struct damacy_executor* executor = NULL; - EXPECT(start_executor(&reader, 2, 2, 8 << 20, &stats, &executor) == 0); + EXPECT(start_executor(&reader, 2, 256, 2, 8 << 20, &stats, &executor) == 0); executor->ops->stats(executor, &stats); uint64_t initial = stats.host_bytes_committed; damacy_executor_destroy(executor); @@ -704,13 +800,14 @@ test_input_memory_budget(void) .max_encoded_chunk_bytes = 2 * CHUNK_BYTES, .max_decoded_chunk_bytes = CHUNK_BYTES, - .max_memory_bytes = 8 << 20 }; + .max_memory_bytes = 8 << 20, + .chunks_per_input_buffer = 256 }; struct damacy_batch_spec output = output_spec(2); EXPECT(damacy_cpu_executor_create(&reader, &config, &executor) == DAMACY_OK); EXPECT(executor->ops->start(executor, &output, &stats) == DAMACY_OK); executor->ops->stats(executor, &stats); uint64_t required = stats.host_bytes_committed; - EXPECT(required == initial + 512 * CHUNK_BYTES); + EXPECT(required == initial + 2 * 256 * CHUNK_BYTES); damacy_executor_destroy(executor); const uint64_t budgets[] = { 1, initial, required - 1, required }; for (unsigned i = 0; i < 4; ++i) { @@ -753,7 +850,9 @@ main(void) RUN(test_plan_memory_retry); RUN(test_read_errors_and_shutdown); RUN(test_read_buffer_capacity); - RUN(test_read_plan_many_workers); + RUN(test_merged_read_limit); + RUN(test_read_plan_large_buffers); RUN(test_input_memory_budget); + RUN(test_input_buffer_validation); return 0; } diff --git a/tests/test_cpu_pipeline.c b/tests/test_cpu_pipeline.c index 92b0150e..d064b6f0 100644 --- a/tests/test_cpu_pipeline.c +++ b/tests/test_cpu_pipeline.c @@ -57,7 +57,8 @@ create_components(struct components* c) &(struct damacy_cpu_config){ .decode_workers = 2, .max_encoded_chunk_bytes = 1024, .max_decoded_chunk_bytes = 1 << 20, - .max_memory_bytes = 32 << 20 }, + .max_memory_bytes = 32 << 20, + .chunks_per_input_buffer = 256 }, &c->executor) == DAMACY_OK); return 0; } From e6b78428277d3a8da5888ee086e91e22e011e5ab Mon Sep 17 00:00:00 2001 From: Nathan Clack Date: Thu, 24 Sep 2026 21:26:31 +0000 Subject: [PATCH 3/5] cpu: log memory budget breakdown --- docs/pipeline.md | 4 +++- src/executor/cpu_executor.c | 24 +++++++++++++++++++++--- 2 files changed, 24 insertions(+), 4 deletions(-) diff --git a/docs/pipeline.md b/docs/pipeline.md index 72bfce8e..58995507 100644 --- a/docs/pipeline.md +++ b/docs/pipeline.md @@ -138,7 +138,9 @@ The two input buffers reserve bookkeeping. With the defaults, 256 chunks and a 4 MiB encoded-chunk bound, that is 2 GiB before decoder workspaces and output buffers. Set the bound to match the largest encoded chunk expected, and include this reserve in -`max_memory_bytes`; insufficient budgets report `BUDGET`. +`max_memory_bytes`. If the budget is too small, starting the pipeline reports +`BUDGET` and logs how many bytes the input buffers, decoder workspaces, and +output buffers need. CPU memory admission includes active read plans, temporary planning scratch, and a conservative allowance for Blosc scratch storage. diff --git a/src/executor/cpu_executor.c b/src/executor/cpu_executor.c index d964c3d5..9cdd4fe2 100644 --- a/src/executor/cpu_executor.c +++ b/src/executor/cpu_executor.c @@ -101,10 +101,17 @@ cpu_read_plan_build(const struct prepared_plan* plan, uint64_t scratch_bytes = (uint64_t)count * (sizeof(struct read_op) + 4 * sizeof(uint32_t)); uint64_t need = bytes + scratch_bytes; - if (need > SIZE_MAX) + if (need > SIZE_MAX || + (need > available && need - available > active_plan_bytes)) { + log_error("read plan for %u chunks needs %llu bytes, but only %llu bytes " + "of max_memory_bytes are left after the CPU executor's buffers", + count, + (unsigned long long)need, + (unsigned long long)(available + active_plan_bytes)); return DAMACY_BUDGET; + } if (need > available) - return need - available <= active_plan_bytes ? DAMACY_AGAIN : DAMACY_BUDGET; + return DAMACY_AGAIN; struct cpu_read_plan reads = { .bytes = bytes }; reads.reads = calloc(1, (size_t)bytes); struct read_op* scratch_reads = calloc(count, sizeof(*scratch_reads)); @@ -451,8 +458,19 @@ cpu_start(struct damacy_executor* base, uint64_t budget = self->config.max_memory_bytes; uint64_t need = fixed + input + decode; if (chunks > SIZE_MAX / self->config.max_encoded_chunk_bytes || - need > budget || bytes > (budget - need) / 2) + need > budget || bytes > (budget - need) / 2) { + log_error("max_memory_bytes=%llu is too small for the CPU executor: " + "input buffers %llu (2 x %u chunks), decode workspace %llu " + "(%u workers), output buffers 2 x %llu, other %llu", + (unsigned long long)budget, + (unsigned long long)input, + (unsigned)chunks, + (unsigned long long)decode, + (unsigned)workers, + (unsigned long long)bytes, + (unsigned long long)fixed); return DAMACY_BUDGET; + } self->committed = need + 2 * bytes; for (unsigned i = 0; i < 2; ++i) { struct damacy_buffer* buffer = calloc(1, sizeof(*buffer)); From 14167be9089cfe7216fecdff10c9e93b552b763e Mon Sep 17 00:00:00 2001 From: Nathan Clack Date: Thu, 24 Sep 2026 21:26:31 +0000 Subject: [PATCH 4/5] bench: fit CPU scenario input buffers --- bench/scenarios/throughput-cpu.json | 5 +++-- dev/cpu-pipeline-validation.md | 4 +++- 2 files changed, 6 insertions(+), 3 deletions(-) diff --git a/bench/scenarios/throughput-cpu.json b/bench/scenarios/throughput-cpu.json index e3c6e127..4417a69d 100644 --- a/bench/scenarios/throughput-cpu.json +++ b/bench/scenarios/throughput-cpu.json @@ -49,7 +49,8 @@ "n_shard_index_cache": 32768, "max_shards_per_sample": 8, "executor": "cpu", - "max_cpu_memory_mb": 4096, - "decode_workers": 8 + "max_cpu_memory_mb": 5120, + "decode_workers": 8, + "chunks_per_input_buffer": 256 } } diff --git a/dev/cpu-pipeline-validation.md b/dev/cpu-pipeline-validation.md index ff434118..62b42050 100644 --- a/dev/cpu-pipeline-validation.md +++ b/dev/cpu-pipeline-validation.md @@ -18,7 +18,9 @@ eight I/O workers, and a 4 GiB executor memory limit. Each worker count was measured once in the same 16-core CPU allocation, in the order 1, 4, 8, 16. Filesystem caches were shared between runs. These are rates after warmup, not cold-storage measurements or training-loop timings. The scenario -is [throughput-cpu.json](../bench/scenarios/throughput-cpu.json). +is [throughput-cpu.json](../bench/scenarios/throughput-cpu.json). It now sets +a 5 GiB limit, because later 256-chunk input buffers reserve 2 GiB with the +default 4 MiB encoded-chunk bound. CUDA comparisons use five warmup batches and thirty measured batches, with a 6 GiB device-memory limit. Baseline and refactor runs alternate on one L40 From 07d503d206478a44bf20a639ade0d3ff99707a99 Mon Sep 17 00:00:00 2001 From: Nathan Clack Date: Thu, 24 Sep 2026 21:57:28 +0000 Subject: [PATCH 5/5] cpu: cap merged reads at decode_workers --- dev/cpu-pipeline.md | 4 ++-- docs/pipeline.md | 7 ++++--- src/executor/cpu_executor.c | 4 +++- tests/test_cpu_executor.c | 36 +++++++++++++++++++++--------------- 4 files changed, 30 insertions(+), 21 deletions(-) diff --git a/dev/cpu-pipeline.md b/dev/cpu-pipeline.md index c7c52cbf..6bd2784e 100644 --- a/dev/cpu-pipeline.md +++ b/dev/cpu-pipeline.md @@ -91,8 +91,8 @@ input groups and two output buffers. Each input group holds at most `chunks_per_input_buffer` chunks and `chunks_per_input_buffer * max_encoded_chunk_bytes` encoded bytes. The setting is at least `decode_workers`, so every worker can get a chunk. Merged reads -contain at most `chunks_per_input_buffer` chunks, and their count respects the -reader capacity. Decoding uses a bounded per-worker output +contain at most `min(decode_workers, chunks_per_input_buffer)` chunks, and their +count respects the reader capacity. Decoding uses a bounded per-worker output buffer and zstd context; memory admission also reserves conservative Blosc workspace. The memory cap covers executor buffers, active read plans, and temporary read-planning diff --git a/docs/pipeline.md b/docs/pipeline.md index 58995507..a2a92a77 100644 --- a/docs/pipeline.md +++ b/docs/pipeline.md @@ -129,9 +129,10 @@ The CPU executor merges adjacent or overlapping encoded ranges within each shard, then interleaves the reads across shards. It decodes each unique source chunk once per batch, including when several output samples use that chunk. Each of its two encoded-input buffers holds `chunks_per_input_buffer` chunks -(default 256). A merged read contains at most that many chunks, so reads merge -even with one decode worker; submission also respects the reader limit. -Decoder workspaces remain per worker. +(default 256). A merged read contains at most +`min(decode_workers, chunks_per_input_buffer)` chunks, so one buffer can hold +several reads at once; with one decode worker, reads do not merge. Submission +also respects the reader limit. Decoder workspaces remain per worker. The two input buffers reserve `2 * chunks_per_input_buffer * max_encoded_chunk_bytes` bytes, plus per-chunk diff --git a/src/executor/cpu_executor.c b/src/executor/cpu_executor.c index 9cdd4fe2..897b5fbf 100644 --- a/src/executor/cpu_executor.c +++ b/src/executor/cpu_executor.c @@ -128,7 +128,9 @@ cpu_read_plan_build(const struct prepared_plan* plan, .file_offset = chunk->offset, .nbytes = chunk->encoded_bytes }; } - uint32_t read_chunks = config->chunks_per_input_buffer; + uint32_t read_chunks = config->decode_workers; + if (read_chunks > config->chunks_per_input_buffer) + read_chunks = config->chunks_per_input_buffer; uint32_t* read_index = indices + 2 * (size_t)count; uint32_t* offset_in_read = indices + 3 * (size_t)count; reads.count = count; diff --git a/tests/test_cpu_executor.c b/tests/test_cpu_executor.c index 0fb28a47..f1f43d8f 100644 --- a/tests/test_cpu_executor.c +++ b/tests/test_cpu_executor.c @@ -656,43 +656,44 @@ test_read_buffer_capacity(void) static int test_merged_read_limit(void) { - const uint32_t buffer_chunks[] = { 64, 256 }; + const uint32_t workers[] = { 1, 2, 32 }; + const uint32_t buffer_chunks = 256; struct test_chunk chunks[513]; for (unsigned i = 0; i < 513; ++i) chunks[i] = (struct test_chunk){ .shard = i % 3, .offset = (170 - i / 3) * CHUNK_BYTES }; - for (unsigned mode = 0; mode < 2; ++mode) { - uint32_t limit = buffer_chunks[mode]; + for (unsigned mode = 0; mode < 3; ++mode) { + uint32_t limit = workers[mode]; struct test_store store; store_init(&store, 512); struct damacy_reader reader = { .store = &store.base, .max_inflight_reads = 512 }; struct damacy_stats stats = { 0 }; struct damacy_executor* executor = NULL; - EXPECT( - start_executor(&reader, 1, limit, 513, 32 << 20, &stats, &executor) == 0); + EXPECT(start_executor( + &reader, limit, buffer_chunks, 513, 32 << 20, &stats, &executor) == + 0); struct prepared_plan* plan = make_plan(chunks, 513); EXPECT(plan); EXPECT(executor->ops->submit(executor, plan, 0) == DAMACY_OK); struct damacy_batch* batch = NULL; EXPECT(finish_batch(executor, &batch) == 0); EXPECT(check_output(batch, &store, chunks, 513) == 0); + EXPECT(stats.chunks_dispatched == 513 && stats.waves_emitted == 3); + EXPECT(store.n_events == 3); + EXPECT(store.events[0].count == buffer_chunks / limit); uint32_t reads_per_shard = (171 + limit - 1) / limit; EXPECT(store.n_reads == 3 * reads_per_shard); - EXPECT(store.n_events == store.n_reads); EXPECT(stats.reads_issued == store.n_reads); - EXPECT(stats.waves_emitted == store.n_reads); for (unsigned i = 0; i < store.n_reads; ++i) { uint32_t first = (i / 3) * limit; uint32_t count = 171 - first; if (count > limit) count = limit; - EXPECT(store.events[i].count == 1); EXPECT(store.records[i].shard == i % 3); EXPECT(store.records[i].offset == (uint64_t)first * CHUNK_BYTES); EXPECT(store.records[i].bytes == count * CHUNK_BYTES); } - EXPECT(stats.chunks_dispatched == 513); EXPECT(stats.io.input_bytes == 513 * CHUNK_BYTES); damacy_batch_release(batch); damacy_executor_destroy(executor); @@ -701,21 +702,26 @@ test_merged_read_limit(void) } static int -test_read_plan_large_buffers(void) +test_read_plan_many_workers(void) { - const uint32_t buffer_chunks[] = { 256, 257, 1024 }; + const struct + { + uint32_t workers; + uint32_t buffer_chunks; + } modes[] = { { 256, 256 }, { 257, 1024 }, { 1024, 1024 } }; struct test_chunk chunks[513]; for (unsigned i = 0; i < 513; ++i) chunks[i] = (struct test_chunk){ .offset = (512 - i) * CHUNK_BYTES }; struct prepared_plan* plan = make_plan(chunks, 513); EXPECT(plan); for (unsigned mode = 0; mode < 3; ++mode) { - uint32_t limit = buffer_chunks[mode]; - struct damacy_cpu_config config = { .decode_workers = 1, + uint32_t limit = modes[mode].workers; + struct damacy_cpu_config config = { .decode_workers = modes[mode].workers, .max_encoded_chunk_bytes = CHUNK_BYTES, .max_decoded_chunk_bytes = CHUNK_BYTES, .max_memory_bytes = 8 << 20, - .chunks_per_input_buffer = limit }; + .chunks_per_input_buffer = + modes[mode].buffer_chunks }; struct cpu_read_plan reads = { 0 }; EXPECT(cpu_read_plan_build(plan, &config, 8 << 20, 0, &reads) == DAMACY_OK); EXPECT(reads.count == (513 + limit - 1) / limit); @@ -851,7 +857,7 @@ main(void) RUN(test_read_errors_and_shutdown); RUN(test_read_buffer_capacity); RUN(test_merged_read_limit); - RUN(test_read_plan_large_buffers); + RUN(test_read_plan_many_workers); RUN(test_input_memory_budget); RUN(test_input_buffer_validation); return 0;