From 280e003475dcd1e7bc826941c629ae8946b6707f Mon Sep 17 00:00:00 2001 From: Nathan Clack Date: Wed, 23 Sep 2026 22:19:35 +0000 Subject: [PATCH 1/3] fix: merge and interleave CPU reads --- dev/cpu-pipeline.md | 21 ++- docs/pipeline.md | 22 ++- src/CMakeLists.txt | 2 +- src/executor/coalesce.c | 14 +- src/executor/coalesce.h | 8 +- src/executor/cpu_executor.c | 186 ++++++++++++++++++---- tests/test_coalesce.c | 34 ++++ tests/test_cpu_executor.c | 305 +++++++++++++++++++++++++++++++++++- 8 files changed, 540 insertions(+), 52 deletions(-) diff --git a/dev/cpu-pipeline.md b/dev/cpu-pipeline.md index 8abd6b8d..ce921d8a 100644 --- a/dev/cpu-pipeline.md +++ b/dev/cpu-pipeline.md @@ -74,8 +74,12 @@ resolution. It preserves the existing GPU decode/assembly path. CUDA chunk layout probes are backend preparation and never enter shared planning. The CPU executor deduplicates decoded source chunks within a batch using the -plan's chunk-use links. It handles raw bytes, zstd, and C-Blosc zstd with no, -byte, or bit shuffle. It copies clipped chunk intersections and casts the +plan's chunk-use links. Before reading, it uses the shared coalescer to merge +adjacent or overlapping ranges within each shard and interleave merged reads +across shards. Each unique chunk retains its offset within a merged read; its +uses still control output assembly. CPU reads use exact encoded byte ranges, +without CUDA's page alignment. It handles raw bytes, zstd, and C-Blosc zstd +with no, byte, or bit shuffle. It copies clipped chunk intersections and casts supported source types to `f32` or `bf16`, including fill-only chunks. The CUDA executor currently retains its per-use decode behavior; cross-sample GPU decode reuse can be optimized separately. @@ -83,11 +87,14 @@ 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 fits the decoder-worker -and reader capacities. Decoding uses a bounded per-worker output buffer and -zstd context; memory admission also reserves conservative Blosc workspace. -The memory cap covers executor buffers, not metadata, reader queues, plan -storage, thread stacks, allocator overhead, or total process RSS. +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 +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. The result handle has its own reference count and owns a reference to its buffer. The executor owns another buffer reference and reuses storage only diff --git a/docs/pipeline.md b/docs/pipeline.md index acdd7db8..f5b1d804 100644 --- a/docs/pipeline.md +++ b/docs/pipeline.md @@ -110,7 +110,7 @@ capacities. Defaults come from the Python value objects. | `PlanLimits.max_chunk_bytes` | Decoded source bytes per chunk, before output conversion. | | `PlanLimits.max_shards_per_sample` | Maximum number of shard files touched by one sample. | | `PlanLimits.max_plan_bytes` | Owned storage per prepared plan. | -| `CpuLimits.max_memory_bytes` | CPU executor's input, decoded, codec-workspace, and two output buffers. | +| `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. | | `FileReader.workers` | Bulk I/O workers, separate from decoding workers. | | `FileReader.max_inflight_reads` | Bulk read capacity; execution respects this bound and retries saturation. | @@ -123,10 +123,22 @@ cover pending samples and the batch being prepared. Queued plans own their metadata and do not pin cache entries. The legacy `Config` adapter retains its older, stricter cache validation. -CPU memory admission includes a conservative allowance for Blosc scratch -storage. `Stats.host_bytes_committed` reports that reservation, which can exceed -the bytes actually touched. It excludes metadata caches, prepared plans, reader -queues, thread stacks, and allocator overhead: it is not a process RSS limit. +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. + +CPU memory admission includes active read plans, temporary planning scratch, +and a conservative allowance for Blosc scratch storage. +`Stats.host_bytes_committed` reports current buffers, read plans, and codec +reservations, which can exceed the bytes actually touched. Temporary planning +scratch is checked against the cap and released before submission returns. +If an active batch's read plan temporarily prevents admission, submission +retries after that batch completes. A plan that cannot fit by itself reports +`BUDGET`. +The cap excludes metadata caches, prepared plans, reader queues, thread stacks, +and allocator overhead: it is not a process RSS limit. Queued plan storage is bounded separately by `prepared_batches * max_plan_bytes`, with up to two accepted plans and one plan being built in addition. Retaining results across repeated pipeline restarts diff --git a/src/CMakeLists.txt b/src/CMakeLists.txt index 27fe67fc..23ee8530 100644 --- a/src/CMakeLists.txt +++ b/src/CMakeLists.txt @@ -279,7 +279,7 @@ if(NOT DAMACY_FUZZ) metadata_store_async lookahead ) add_src_lib(cpu_executor SOURCES executor/cpu_executor.c - LINKS pipeline_components prepared_plan threadpool damacy_stats + LINKS pipeline_components prepared_plan dispatch_utils threadpool damacy_stats PkgConfig::ZSTD PkgConfig::BLOSC ) add_src_lib(damacy diff --git a/src/executor/coalesce.c b/src/executor/coalesce.c index 1cb70981..0f2337d7 100644 --- a/src/executor/coalesce.c +++ b/src/executor/coalesce.c @@ -2,6 +2,14 @@ #include "executor/read_op_sort.h" +#include + +static int +same_path(const char* a, const char* b) +{ + return a == b || !strcmp(a, b); +} + enum damacy_status coalesce_chunks(struct dispatch_output* out, uint64_t read_op_max_bytes, @@ -28,7 +36,7 @@ coalesce_chunks(struct dispatch_output* out, for (uint32_t i = 0; i < n; ++i) remap[i] = UINT32_MAX; - // Partition: real (path interned, nbytes > 0) vs fill placeholders. + // Partition: real (path present, nbytes > 0) vs fill placeholders. uint32_t n_io = 0; for (uint32_t i = 0; i < n; ++i) { struct read_op* r = &out->read_ops[i]; @@ -51,7 +59,7 @@ coalesce_chunks(struct dispatch_output* out, int fusable = 0; if (leader_old != UINT32_MAX) { struct read_op* leader = &out->read_ops[leader_old]; - if (curr->shard_path == leader->shard_path && + if (same_path(curr->shard_path, leader->shard_path) && curr->file_offset >= leader->file_offset && curr->file_offset <= leader_end && leader_chunks < max_chunks_per_wave) { @@ -91,7 +99,7 @@ coalesce_chunks(struct dispatch_output* out, uint32_t* starts = perm; // perm is dead after the fuse loop uint32_t n_runs = 0; for (uint32_t k = 0; k < n_real; ++k) - if (k == 0 || tmp[k].shard_path != tmp[k - 1].shard_path) + if (k == 0 || !same_path(tmp[k].shard_path, tmp[k - 1].shard_path)) starts[n_runs++] = k; uint32_t outk = 0; for (uint32_t round = 0; outk < n_real; ++round) { diff --git a/src/executor/coalesce.h b/src/executor/coalesce.h index 118e3c66..e5b56dad 100644 --- a/src/executor/coalesce.h +++ b/src/executor/coalesce.h @@ -1,16 +1,16 @@ // Coalesce step of the IO planning pipeline: // filter (emit) → sort → fuse-with-cap → interleave → group-by-read. // -// Operates in place on dispatch_output: sorts the per-chunk -// page-aligned read windows by (shard_path, file_offset), then +// Operates in place on dispatch_output: sorts per-chunk read windows +// by (shard_path, file_offset), then // greedily fuses adjacent windows in the same shard into one // read_op, bounded by read_op_max_bytes. Fused ops are emitted // interleaved across shards (per-shard offset order kept) because // simultaneous reads into one file serialize on network // filesystems. chunk_plan.read_op_idx and // chunk_plan.offset_in_read are rewritten to point at the surviving -// read_op. Fill chunk_plans (path empty, nbytes == 0) keep their -// 1:1 placeholder read_ops untouched. +// read_op. Equal shard paths need not share a pointer. Fill chunk_plans +// (path empty, nbytes == 0) keep their 1:1 placeholder read_ops untouched. // // Populates dispatch_output.n_chunks_to_load and n_loads_issued as // part of the same pass. diff --git a/src/executor/cpu_executor.c b/src/executor/cpu_executor.c index e44df4b9..d6f48bf2 100644 --- a/src/executor/cpu_executor.c +++ b/src/executor/cpu_executor.c @@ -6,6 +6,7 @@ #include "damacy_config.h" #include "damacy_stats.h" +#include "executor/coalesce.h" #include "threadpool/threadpool.h" #include @@ -19,11 +20,27 @@ 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; struct damacy_buffer* buffer; struct prepared_plan* plan; + struct cpu_read_plan read_plan; uint64_t batch_id; uint32_t dispatched; uint32_t remaining; @@ -36,17 +53,23 @@ struct cpu_worker ZSTD_DCtx* zstd; }; +struct cpu_chunk_input +{ + const struct plan_chunk* chunk; + const void* input; +}; + struct cpu_wave { struct store_event event; struct store_read* reads; + struct cpu_chunk_input* chunks; void* input; enum damacy_status* results; float* decode_ms; float* assemble_ms; uint64_t* output_bytes; struct platform_clock clock; - uint32_t first_chunk; uint32_t count; uint32_t slot; uint64_t input_bytes; @@ -70,6 +93,83 @@ struct cpu_executor uint64_t committed; }; +static void +cpu_read_plan_destroy(struct cpu_read_plan* plan) +{ + free(plan->reads); + *plan = (struct cpu_read_plan){ 0 }; +} + +static 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) +{ + uint32_t count = plan->n_chunks; + uint64_t bytes = + (uint64_t)count * + (sizeof(struct read_op) + sizeof(struct cpu_chunk_read) + sizeof(uint32_t)); + uint64_t scratch_bytes = + (uint64_t)count * + (sizeof(struct chunk_plan) + sizeof(struct read_op) + 4 * sizeof(uint32_t)); + uint64_t need = bytes + scratch_bytes; + if (need > SIZE_MAX) + return DAMACY_BUDGET; + if (need > available) + return need - available <= active_plan_bytes ? DAMACY_AGAIN : DAMACY_BUDGET; + struct cpu_read_plan reads = { .bytes = bytes }; + reads.reads = calloc(1, (size_t)bytes); + struct chunk_plan* chunks = calloc(count, sizeof(*chunks)); + struct read_op* scratch_reads = calloc(count, sizeof(*scratch_reads)); + uint32_t* indices = calloc((size_t)count * 4, sizeof(*indices)); + enum damacy_status status = DAMACY_OOM; + if (!reads.reads || !chunks || !scratch_reads || !indices) + goto Done; + reads.chunks = (void*)(reads.reads + count); + reads.first_chunks = (void*)(reads.chunks + count); + for (uint32_t i = 0; i < count; ++i) { + const struct plan_chunk* chunk = &plan->chunks[i]; + chunks[i] = + (struct chunk_plan){ .read_op_idx = i, .is_fill = chunk->missing }; + if (!chunk->missing) + reads.reads[i] = (struct read_op){ .shard_path = chunk->path, + .file_offset = chunk->offset, + .nbytes = chunk->encoded_bytes }; + } + struct dispatch_output dispatch = { .read_ops = reads.reads, + .n_read_ops = count, + .chunk_plans = chunks, + .n_chunk_plans = count }; + status = coalesce_chunks(&dispatch, + (uint64_t)config->decode_workers * + config->max_encoded_chunk_bytes, + config->decode_workers, + indices, + scratch_reads); + if (status != DAMACY_OK) + goto Done; + reads.count = dispatch.n_read_ops; + for (uint32_t i = 0; i < reads.count; ++i) + reads.first_chunks[i] = UINT32_MAX; + for (uint32_t i = count; i-- > 0;) { + uint32_t read = chunks[i].read_op_idx; + reads.chunks[i] = + (struct cpu_chunk_read){ .offset = chunks[i].offset_in_read, + .next = reads.first_chunks[read] }; + reads.first_chunks[read] = i; + } + *out = reads; +Done: + free(chunks); + free(scratch_reads); + free(indices); + if (status != DAMACY_OK) + cpu_read_plan_destroy(&reads); + return status; +} + static float half_to_float(uint16_t half) { @@ -277,11 +377,9 @@ decode_one(size_t index, int tid, void* arg) struct cpu_executor* self = arg; struct cpu_wave* wave = self->decoding; struct cpu_slot* slot = &self->slots[wave->slot]; - const struct plan_chunk* chunk = - &slot->plan->chunks[wave->first_chunk + index]; + const struct plan_chunk* chunk = wave->chunks[index].chunk; const struct zarr_metadata* meta = &slot->plan->arrays[chunk->array].metadata; - const void* input = - (const char*)wave->input + index * self->config.max_encoded_chunk_bytes; + const void* input = wave->chunks[index].input; const void* decoded; struct platform_clock clock = { 0 }; platform_toc(&clock); @@ -313,6 +411,7 @@ cpu_stop(struct damacy_executor* base) store_event_wait(self->reader->store, wave->event); free(wave->input); free(wave->reads); + free(wave->chunks); free(wave->results); free(wave->decode_ms); free(wave->assemble_ms); @@ -321,6 +420,7 @@ cpu_stop(struct damacy_executor* base) } for (unsigned i = 0; i < 2; ++i) { prepared_plan_destroy(self->slots[i].plan); + cpu_read_plan_destroy(&self->slots[i].read_plan); buffer_release(self->slots[i].buffer); self->slots[i] = (struct cpu_slot){ 0 }; } @@ -358,8 +458,8 @@ cpu_start(struct damacy_executor* base, 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(enum damacy_status) + - 2 * sizeof(float) + sizeof(uint64_t)); + 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) @@ -384,12 +484,13 @@ cpu_start(struct damacy_executor* base, 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)); - if (!wave->input || !wave->reads || !wave->results || !wave->decode_ms || - !wave->assemble_ms || !wave->output_bytes) + if (!wave->input || !wave->reads || !wave->chunks || !wave->results || + !wave->decode_ms || !wave->assemble_ms || !wave->output_bytes) goto Fail; } self->workers = calloc((size_t)workers, sizeof(*self->workers)); @@ -428,6 +529,8 @@ cpu_submit(struct damacy_executor* base, } if (slot < 0) return DAMACY_AGAIN; + if (!plan || !plan->n_chunks) + return DAMACY_INVAL; for (uint32_t i = 0; i < plan->n_chunks; ++i) { const struct plan_chunk* chunk = &plan->chunks[i]; if (chunk->encoded_bytes > self->config.max_encoded_chunk_bytes || @@ -440,6 +543,15 @@ cpu_submit(struct damacy_executor* base, return DAMACY_DECODE; } struct cpu_slot* target = &self->slots[slot]; + enum damacy_status status = cpu_read_plan_build( + plan, + &self->config, + self->config.max_memory_bytes - self->committed, + self->slots[0].read_plan.bytes + self->slots[1].read_plan.bytes, + &target->read_plan); + if (status != DAMACY_OK) + return status; + self->committed += target->read_plan.bytes; target->plan = plan; target->batch_id = batch_id; target->dispatched = 0; @@ -455,32 +567,44 @@ cpu_dispatch(struct cpu_executor* self, struct cpu_wave* wave, int* changed) for (unsigned i = 0; i < 2; ++i) { const struct cpu_slot* slot = &self->slots[i]; if (slot->state == CPU_RENDERING && - slot->dispatched < slot->plan->n_chunks && + slot->dispatched < slot->read_plan.count && (slot_index < 0 || slot->batch_id < self->slots[slot_index].batch_id)) slot_index = (int)i; } if (slot_index < 0) return DAMACY_OK; struct cpu_slot* slot = &self->slots[slot_index]; - uint32_t count = slot->plan->n_chunks - slot->dispatched; - if (count > self->config.decode_workers) - count = self->config.decode_workers; - if (count > self->reader->max_inflight_reads) - count = self->reader->max_inflight_reads; + const struct cpu_read_plan* plan = &slot->read_plan; + uint32_t next_read = slot->dispatched; + uint32_t count = 0; uint32_t n_reads = 0; uint64_t bytes = 0; - for (uint32_t i = 0; i < count; ++i) { - const struct plan_chunk* chunk = &slot->plan->chunks[slot->dispatched + i]; - if (chunk->missing) - continue; - wave->reads[n_reads++] = - (struct store_read){ .key = chunk->path, - .offset = chunk->offset, - .len = chunk->encoded_bytes, - .dst = - (char*)wave->input + - (size_t)i * self->config.max_encoded_chunk_bytes }; - bytes += chunk->encoded_bytes; + uint64_t capacity = (uint64_t)self->config.decode_workers * + 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 || + 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; + c = plan->chunks[c].next) + wave->chunks[count++] = + (struct cpu_chunk_input){ .chunk = &slot->plan->chunks[c], + .input = (char*)wave->input + bytes + + plan->chunks[c].offset }; + if (read->nbytes) + wave->reads[n_reads++] = + (struct store_read){ .key = read->shard_path, + .offset = read->file_offset, + .len = read->nbytes, + .dst = (char*)wave->input + bytes }; + bytes += read->nbytes; + ++next_read; } platform_toc(&wave->clock); struct store_submit_result result = @@ -490,10 +614,9 @@ cpu_dispatch(struct cpu_executor* self, struct cpu_wave* wave, int* changed) wave->event = result.event; wave->active = 1; wave->slot = (uint32_t)slot_index; - wave->first_chunk = slot->dispatched; wave->count = count; wave->input_bytes = bytes; - slot->dispatched += count; + slot->dispatched = next_read; self->stats->reads_issued += n_reads; self->stats->chunks_dispatched += count; ++self->stats->waves_emitted; @@ -540,8 +663,7 @@ cpu_step(struct damacy_executor* base, int* changed) for (uint32_t j = 0; j < wave->count; ++j) { if (wave->results[j] != DAMACY_OK) return wave->results[j]; - const struct plan_chunk* chunk = - &slot->plan->chunks[wave->first_chunk + j]; + const struct plan_chunk* chunk = wave->chunks[j].chunk; metric_record(&self->stats->decode, wave->decode_ms[j], chunk->encoded_bytes, @@ -555,6 +677,8 @@ cpu_step(struct damacy_executor* base, int* changed) wave->active = 0; if (!slot->remaining) { slot->state = CPU_READY; + self->committed -= slot->read_plan.bytes; + cpu_read_plan_destroy(&slot->read_plan); prepared_plan_destroy(slot->plan); slot->plan = NULL; } diff --git a/tests/test_coalesce.c b/tests/test_coalesce.c index 70af6d19..a06a1605 100644 --- a/tests/test_coalesce.c +++ b/tests/test_coalesce.c @@ -422,6 +422,39 @@ test_fills_passthrough(void) return 0; } +static int +test_equal_path_copies(void) +{ + char paths[6][8] = { "shard/0", "shard/0", "shard/0", + "shard/1", "shard/1", "shard/1" }; + const uint64_t offsets[] = { 0, 4, 16, 0, 4, 16 }; + const uint32_t expected_reads[] = { 0, 0, 2, 1, 1, 3 }; + struct read_op reads[6] = { 0 }; + struct chunk_plan chunks[6] = { 0 }; + for (uint32_t i = 0; i < 6; ++i) { + reads[i] = (struct read_op){ .shard_path = paths[i], + .file_offset = offsets[i], + .nbytes = 4 }; + chunks[i].read_op_idx = i; + } + struct dispatch_output out = { .read_ops = reads, + .n_read_ops = 6, + .chunk_plans = chunks, + .n_chunk_plans = 6 }; + EXPECT(run_coalesce(&out, 8, 6) == DAMACY_OK); + EXPECT(out.n_read_ops == 4); + for (uint32_t i = 0; i < 4; ++i) { + EXPECT(strcmp(reads[i].shard_path, i % 2 ? "shard/1" : "shard/0") == 0); + EXPECT(reads[i].file_offset == (i < 2 ? 0u : 16u)); + EXPECT(reads[i].nbytes == (i < 2 ? 8u : 4u)); + } + for (uint32_t i = 0; i < 6; ++i) { + EXPECT(chunks[i].read_op_idx == expected_reads[i]); + EXPECT(chunks[i].offset_in_read == (i % 3 == 1 ? 4u : 0u)); + } + return 0; +} + int main(void) { @@ -437,5 +470,6 @@ main(void) RUN(test_round_robin_interleave); RUN(test_chunk_count_cap); RUN(test_fills_passthrough); + RUN(test_equal_path_copies); return 0; } diff --git a/tests/test_cpu_executor.c b/tests/test_cpu_executor.c index 01ddff14..75c4517c 100644 --- a/tests/test_cpu_executor.c +++ b/tests/test_cpu_executor.c @@ -229,6 +229,302 @@ start_executor(struct damacy_reader* reader, return 0; } +static int +finish_batch(struct damacy_executor* executor, struct damacy_batch** out) +{ + for (unsigned i = 0; i < 2 * MAX_CHUNKS; ++i) { + int changed = 0; + EXPECT(executor->ops->step(executor, &changed) == DAMACY_OK); + enum damacy_status status = executor->ops->take(executor, out); + if (status == DAMACY_OK) + return 0; + EXPECT(status == DAMACY_AGAIN && changed); + } + return 1; +} + +static int +check_output(struct damacy_batch* batch, + const struct test_store* store, + const struct test_chunk* chunks, + uint32_t count) +{ + struct damacy_batch_info info; + damacy_batch_info(batch, &info); + EXPECT(info.device_type == DAMACY_DEVICE_CPU); + EXPECT(info.rank == 2 && info.shape[0] == 2 && info.shape[1] == count * 4); + const float* data = info.data; + for (unsigned sample = 0; sample < 2; ++sample) + for (uint32_t c = 0; c < count; ++c) + for (unsigned i = 0; i < 4; ++i) { + float expected = + chunks[c].missing + ? 999 + : store->data[chunks[c].shard][chunks[c].offset / 2 + i]; + EXPECT(data[(sample * count + c) * 4 + i] == expected); + } + return 0; +} + +static int +test_merge_and_shard_order(void) +{ + struct test_chunk chunks[12]; + for (unsigned shard = 0; shard < 3; ++shard) + for (unsigned i = 0; i < 4; ++i) + chunks[shard * 4 + i] = + (struct test_chunk){ .shard = shard, + .offset = (i < 2 ? 40 : 8) - (i % 2) * 8 }; + struct test_store store; + store_init(&store, 4); + struct damacy_reader reader = { .store = &store.base, + .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); + struct prepared_plan* plan = make_plan(chunks, 12); + EXPECT(plan); + EXPECT(plan->chunks[0].path != plan->chunks[1].path); + 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, 12) == 0); + EXPECT(store.n_reads == 6); + for (unsigned i = 0; i < store.n_reads; ++i) { + EXPECT(store.records[i].shard == i % 3); + EXPECT(store.records[i].offset == (i < 3 ? 0 : 32)); + EXPECT(store.records[i].bytes == 2 * CHUNK_BYTES); + } + EXPECT(stats.reads_issued == 6 && stats.chunks_dispatched == 12); + EXPECT(stats.decode.count == 12); + EXPECT(stats.decode.input_bytes == 12 * CHUNK_BYTES); + EXPECT(stats.decode.output_bytes == 12 * CHUNK_BYTES); + EXPECT(stats.io.input_bytes == 12 * CHUNK_BYTES); + damacy_batch_release(batch); + damacy_executor_destroy(executor); + return 0; +} + +static int +test_reader_limit_and_retry(void) +{ + const struct test_chunk chunks[] = { + { .offset = 0 }, { .offset = 8 }, { .offset = 16 }, { .offset = 24 } + }; + struct test_store store; + store_init(&store, 1); + struct damacy_reader reader = { .store = &store.base, + .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); + struct prepared_plan* plan = make_plan(chunks, 4); + 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, 4) == 0); + EXPECT(store.rejected > 0 && store.pending == 0 && store.n_reads == 2); + EXPECT(store.records[0].offset == 0 && store.records[0].bytes == 16); + EXPECT(store.records[1].offset == 16 && store.records[1].bytes == 16); + EXPECT(stats.chunks_dispatched == 4 && stats.reads_issued == 2); + damacy_batch_release(batch); + damacy_executor_destroy(executor); + return 0; +} + +static int +test_overlapping_ranges(void) +{ + const struct test_chunk chunks[] = { { .offset = 4 }, { .offset = 0 } }; + 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); + struct prepared_plan* plan = make_plan(chunks, 2); + 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, 2) == 0); + EXPECT(store.n_reads == 1 && store.records[0].offset == 0); + EXPECT(store.records[0].bytes == 12); + EXPECT(stats.decode.input_bytes == 16 && stats.io.input_bytes == 12); + damacy_batch_release(batch); + damacy_executor_destroy(executor); + return 0; +} + +static int +test_single_worker_and_fills(void) +{ + for (unsigned all_missing = 0; all_missing < 2; ++all_missing) { + const struct test_chunk chunks[] = { + { .missing = 1 }, + { .shard = 2, .offset = 24, .missing = (uint8_t)all_missing }, + { .missing = 1 }, + { .shard = 0, .offset = 8, .missing = (uint8_t)all_missing }, + }; + struct test_store store; + store_init(&store, 1); + struct damacy_reader reader = { .store = &store.base, + .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); + struct prepared_plan* plan = make_plan(chunks, 4); + 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, 4) == 0); + EXPECT(store.n_reads == (all_missing ? 0u : 2u)); + EXPECT(stats.chunks_dispatched == 4); + damacy_batch_release(batch); + damacy_executor_destroy(executor); + } + return 0; +} + +static int +test_batch_order_and_retained_output(void) +{ + const struct test_chunk chunks[] = { { .offset = 8 }, { .offset = 0 } }; + struct test_store store; + store_init(&store, 2); + store.held_event = 1; + 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); + for (unsigned i = 0; i < 2; ++i) { + struct prepared_plan* plan = make_plan(chunks, 2); + EXPECT(plan); + EXPECT(executor->ops->submit(executor, plan, i) == DAMACY_OK); + } + int changed = 0; + EXPECT(executor->ops->step(executor, &changed) == DAMACY_OK); + EXPECT(store.n_events == 2 && store.pending == 1); + struct damacy_batch* first = NULL; + EXPECT(executor->ops->take(executor, &first) == DAMACY_AGAIN); + store.held_event = 0; + EXPECT(finish_batch(executor, &first) == 0); + struct damacy_batch* second = NULL; + EXPECT(executor->ops->take(executor, &second) == DAMACY_OK); + struct damacy_batch_info a, b; + damacy_batch_info(first, &a); + damacy_batch_info(second, &b); + EXPECT(a.batch_id == 0 && b.batch_id == 1 && a.data != b.data); + struct prepared_plan* third = make_plan(chunks, 2); + EXPECT(third); + EXPECT(executor->ops->submit(executor, third, 2) == DAMACY_AGAIN); + damacy_batch_retain(first); + damacy_batch_release(first); + damacy_batch_release(second); + EXPECT(executor->ops->step(executor, &changed) == DAMACY_OK); + EXPECT(executor->ops->submit(executor, third, 2) == DAMACY_OK); + struct damacy_batch* last = NULL; + EXPECT(finish_batch(executor, &last) == 0); + EXPECT(check_output(last, &store, chunks, 2) == 0); + damacy_batch_release(last); + damacy_executor_destroy(executor); + EXPECT(check_output(first, &store, chunks, 2) == 0); + damacy_batch_release(first); + return 0; +} + +static int +test_plan_memory_budget(void) +{ + const struct test_chunk chunks[] = { { .offset = 0 }, { .offset = 8 } }; + 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; + struct prepared_plan* plan = make_plan(chunks, 2); + EXPECT(plan); + EXPECT(executor->ops->submit(executor, plan, 0) == DAMACY_OK); + executor->ops->stats(executor, &stats); + uint64_t with_plan = stats.host_bytes_committed; + EXPECT(with_plan > initial && with_plan <= (8 << 20)); + struct damacy_batch* batch = NULL; + EXPECT(finish_batch(executor, &batch) == 0); + executor->ops->stats(executor, &stats); + EXPECT(stats.host_bytes_committed == initial); + damacy_batch_release(batch); + damacy_executor_destroy(executor); + executor = NULL; + EXPECT(start_executor(&reader, 2, 2, with_plan, &stats, &executor) == 0); + plan = make_plan(chunks, 2); + EXPECT(plan); + EXPECT(executor->ops->submit(executor, plan, 0) == DAMACY_BUDGET); + executor->ops->stats(executor, &stats); + EXPECT(stats.host_bytes_committed == initial); + prepared_plan_destroy(plan); + damacy_executor_destroy(executor); + return 0; +} + +static int +test_plan_memory_retry(void) +{ + const struct test_chunk chunks[] = { { .offset = 0 }, { .offset = 8 } }; + 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; + struct prepared_plan* plan = make_plan(chunks, 2); + EXPECT(plan); + EXPECT(executor->ops->submit(executor, plan, 0) == DAMACY_OK); + executor->ops->stats(executor, &stats); + uint64_t plan_bytes = stats.host_bytes_committed - initial; + EXPECT(plan_bytes > 0); + damacy_executor_destroy(executor); + for (unsigned i = 1; i <= 16; ++i) { + executor = NULL; + EXPECT(start_executor( + &reader, 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); + if (status == DAMACY_BUDGET) { + prepared_plan_destroy(plan); + damacy_executor_destroy(executor); + continue; + } + EXPECT(status == DAMACY_OK); + struct prepared_plan* next = make_plan(chunks, 2); + EXPECT(next); + EXPECT(executor->ops->submit(executor, next, 1) == DAMACY_AGAIN); + struct damacy_batch* first = NULL; + EXPECT(finish_batch(executor, &first) == 0); + EXPECT(executor->ops->submit(executor, next, 1) == DAMACY_OK); + struct damacy_batch* second = NULL; + EXPECT(finish_batch(executor, &second) == 0); + EXPECT(check_output(first, &store, chunks, 2) == 0); + EXPECT(check_output(second, &store, chunks, 2) == 0); + damacy_batch_release(first); + damacy_batch_release(second); + damacy_executor_destroy(executor); + return 0; + } + return 1; +} + static int test_read_errors_and_shutdown(void) { @@ -237,7 +533,7 @@ test_read_errors_and_shutdown(void) }; for (unsigned failure = 0; failure < 3; ++failure) { struct test_store store; - store_init(&store, 4); + store_init(&store, 2); if (failure == 0) store.submit_status = DAMACY_IO; else if (failure == 1) @@ -270,6 +566,13 @@ test_read_errors_and_shutdown(void) int main(void) { + RUN(test_merge_and_shard_order); + RUN(test_reader_limit_and_retry); + RUN(test_overlapping_ranges); + RUN(test_single_worker_and_fills); + RUN(test_batch_order_and_retained_output); + RUN(test_plan_memory_budget); + RUN(test_plan_memory_retry); RUN(test_read_errors_and_shutdown); return 0; } From 13e841fdd77d78d7bf9582363bf837a5f2b5a7f6 Mon Sep 17 00:00:00 2001 From: Nathan Clack Date: Thu, 24 Sep 2026 21:13:11 +0000 Subject: [PATCH 2/3] cpu: plan reads without chunk_plan --- src/executor/coalesce.c | 100 ++++++++++++++++++++++++------------ src/executor/coalesce.h | 15 ++++++ src/executor/cpu_executor.c | 36 ++++++------- 3 files changed, 97 insertions(+), 54 deletions(-) diff --git a/src/executor/coalesce.c b/src/executor/coalesce.c index 0f2337d7..7a0d0ffa 100644 --- a/src/executor/coalesce.c +++ b/src/executor/coalesce.c @@ -11,40 +11,38 @@ same_path(const char* a, const char* b) } enum damacy_status -coalesce_chunks(struct dispatch_output* out, - uint64_t read_op_max_bytes, - uint32_t max_chunks_per_wave, - uint32_t* u32_scratch, - struct read_op* read_op_scratch) +coalesce_reads(struct read_op* reads, + uint32_t* n_reads, + uint64_t read_op_max_bytes, + uint32_t max_chunks_per_wave, + uint32_t* read_index, + uint32_t* offset_in_read, + uint32_t* u32_scratch, + struct read_op* read_op_scratch) { - if (!out || !out->read_ops || !out->chunk_plans) + if (!reads || !n_reads) return DAMACY_INVAL; - out->n_chunks_to_load = 0; - out->n_loads_issued = 0; - - uint32_t n = out->n_read_ops; + uint32_t n = *n_reads; if (n == 0) return DAMACY_OK; - if (!u32_scratch || !read_op_scratch) + if (!read_index || !offset_in_read || !u32_scratch || !read_op_scratch) return DAMACY_INVAL; uint32_t* perm = u32_scratch; - uint32_t* remap = u32_scratch + n; - uint32_t* offset_shift = u32_scratch + 2u * n; struct read_op* tmp = read_op_scratch; for (uint32_t i = 0; i < n; ++i) - remap[i] = UINT32_MAX; + read_index[i] = UINT32_MAX; // Partition: real (path present, nbytes > 0) vs fill placeholders. uint32_t n_io = 0; for (uint32_t i = 0; i < n; ++i) { - struct read_op* r = &out->read_ops[i]; + struct read_op* r = &reads[i]; if (r->nbytes != 0 && r->shard_path) perm[n_io++] = i; } - read_op_perm_sort(out->read_ops, perm, n_io); + read_op_perm_sort(reads, perm, n_io); // leader_chunks = post-coalesce group size; capped so each group // fits one wave's chunk intake (planner is 1 read_op per chunk). @@ -54,11 +52,11 @@ coalesce_chunks(struct dispatch_output* out, uint32_t leader_chunks = 0; for (uint32_t k = 0; k < n_io; ++k) { uint32_t e = perm[k]; - struct read_op* curr = &out->read_ops[e]; + struct read_op* curr = &reads[e]; uint64_t curr_end = (uint64_t)curr->file_offset + curr->nbytes; int fusable = 0; if (leader_old != UINT32_MAX) { - struct read_op* leader = &out->read_ops[leader_old]; + struct read_op* leader = &reads[leader_old]; if (same_path(curr->shard_path, leader->shard_path) && curr->file_offset >= leader->file_offset && curr->file_offset <= leader_end && @@ -70,18 +68,18 @@ coalesce_chunks(struct dispatch_output* out, } } if (fusable) { - uint32_t new_idx = remap[leader_old]; + uint32_t new_idx = read_index[leader_old]; uint64_t fused_end = curr_end > leader_end ? curr_end : leader_end; tmp[new_idx].nbytes = (uint32_t)(fused_end - tmp[new_idx].file_offset); - remap[e] = new_idx; - offset_shift[e] = + read_index[e] = new_idx; + offset_in_read[e] = (uint32_t)(curr->file_offset - tmp[new_idx].file_offset); leader_end = fused_end; leader_chunks++; } else { tmp[write] = *curr; - remap[e] = write; - offset_shift[e] = 0; + read_index[e] = write; + offset_in_read[e] = 0; leader_old = e; leader_end = curr_end; leader_chunks = 1; @@ -94,7 +92,7 @@ coalesce_chunks(struct dispatch_output* out, // Simultaneous reads into one file serialize on network filesystems // (~2x slower bulk reads), so spread the emitted order across shards. // pos[k] = output slot of tmp[k]; fills below get identity slots. - uint32_t* pos = u32_scratch + 3u * n; + uint32_t* pos = u32_scratch + n; { uint32_t* starts = perm; // perm is dead after the fuse loop uint32_t n_runs = 0; @@ -113,29 +111,63 @@ coalesce_chunks(struct dispatch_output* out, } // Fill placeholders (path empty / nbytes == 0): keep 1:1, append at - // the end of the output. chunk_plans referencing them still find - // their entry via remap. + // the end of the output. for (uint32_t i = 0; i < n; ++i) { - if (remap[i] != UINT32_MAX) + if (read_index[i] != UINT32_MAX) continue; - tmp[write] = out->read_ops[i]; - remap[i] = write; - offset_shift[i] = 0; + tmp[write] = reads[i]; + read_index[i] = write; + offset_in_read[i] = 0; pos[write] = write; write++; } for (uint32_t k = 0; k < write; ++k) - out->read_ops[pos[k]] = tmp[k]; - out->n_read_ops = write; + reads[pos[k]] = tmp[k]; + for (uint32_t i = 0; i < n; ++i) + read_index[i] = pos[read_index[i]]; + *n_reads = write; + return DAMACY_OK; +} + +enum damacy_status +coalesce_chunks(struct dispatch_output* out, + uint64_t read_op_max_bytes, + uint32_t max_chunks_per_wave, + uint32_t* u32_scratch, + struct read_op* read_op_scratch) +{ + if (!out || !out->read_ops || !out->chunk_plans) + return DAMACY_INVAL; + out->n_chunks_to_load = 0; + out->n_loads_issued = 0; + + uint32_t n = out->n_read_ops; + if (n == 0) + return DAMACY_OK; + if (!u32_scratch || !read_op_scratch) + return DAMACY_INVAL; + + uint32_t* read_index = u32_scratch + 2u * n; + uint32_t* offset_in_read = u32_scratch + 3u * n; + enum damacy_status status = coalesce_reads(out->read_ops, + &out->n_read_ops, + read_op_max_bytes, + max_chunks_per_wave, + read_index, + offset_in_read, + u32_scratch, + read_op_scratch); + if (status != DAMACY_OK) + return status; for (uint32_t i = 0; i < out->n_chunk_plans; ++i) { struct chunk_plan* cp = &out->chunk_plans[i]; uint32_t old = cp->read_op_idx; - cp->read_op_idx = pos[remap[old]]; + cp->read_op_idx = read_index[old]; if (cp->is_fill) continue; - uint64_t sum = (uint64_t)cp->offset_in_read + (uint64_t)offset_shift[old]; + uint64_t sum = (uint64_t)cp->offset_in_read + (uint64_t)offset_in_read[old]; if (sum > UINT32_MAX) return DAMACY_INVAL; cp->offset_in_read = (uint32_t)sum; diff --git a/src/executor/coalesce.h b/src/executor/coalesce.h index e5b56dad..23490549 100644 --- a/src/executor/coalesce.h +++ b/src/executor/coalesce.h @@ -31,6 +31,21 @@ extern "C" { #endif + // Merges reads in place and updates *n_reads. read_index[i] and + // offset_in_read[i] locate original read i in the merged list. + // Scratch requirements, with n the original *n_reads: + // read_index, offset_in_read: >= n uint32_t slots each + // u32_scratch: >= 2 * n uint32_t slots + // read_op_scratch: >= n slots + enum damacy_status coalesce_reads(struct read_op* reads, + uint32_t* n_reads, + uint64_t read_op_max_bytes, + uint32_t max_chunks_per_wave, + uint32_t* read_index, + uint32_t* offset_in_read, + uint32_t* u32_scratch, + struct read_op* read_op_scratch); + // Scratch requirements: // u32_scratch: >= 4 * out->n_read_ops uint32_t slots // read_op_scratch: >= out->n_read_ops slots diff --git a/src/executor/cpu_executor.c b/src/executor/cpu_executor.c index d6f48bf2..42efb26a 100644 --- a/src/executor/cpu_executor.c +++ b/src/executor/cpu_executor.c @@ -112,8 +112,7 @@ cpu_read_plan_build(const struct prepared_plan* plan, (uint64_t)count * (sizeof(struct read_op) + sizeof(struct cpu_chunk_read) + sizeof(uint32_t)); uint64_t scratch_bytes = - (uint64_t)count * - (sizeof(struct chunk_plan) + sizeof(struct read_op) + 4 * sizeof(uint32_t)); + (uint64_t)count * (sizeof(struct read_op) + 4 * sizeof(uint32_t)); uint64_t need = bytes + scratch_bytes; if (need > SIZE_MAX) return DAMACY_BUDGET; @@ -121,48 +120,45 @@ cpu_read_plan_build(const struct prepared_plan* plan, return need - available <= active_plan_bytes ? DAMACY_AGAIN : DAMACY_BUDGET; struct cpu_read_plan reads = { .bytes = bytes }; reads.reads = calloc(1, (size_t)bytes); - struct chunk_plan* chunks = calloc(count, sizeof(*chunks)); struct read_op* scratch_reads = calloc(count, sizeof(*scratch_reads)); uint32_t* indices = calloc((size_t)count * 4, sizeof(*indices)); enum damacy_status status = DAMACY_OOM; - if (!reads.reads || !chunks || !scratch_reads || !indices) + if (!reads.reads || !scratch_reads || !indices) goto Done; reads.chunks = (void*)(reads.reads + count); reads.first_chunks = (void*)(reads.chunks + count); for (uint32_t i = 0; i < count; ++i) { const struct plan_chunk* chunk = &plan->chunks[i]; - chunks[i] = - (struct chunk_plan){ .read_op_idx = i, .is_fill = chunk->missing }; if (!chunk->missing) reads.reads[i] = (struct read_op){ .shard_path = chunk->path, .file_offset = chunk->offset, .nbytes = chunk->encoded_bytes }; } - struct dispatch_output dispatch = { .read_ops = reads.reads, - .n_read_ops = count, - .chunk_plans = chunks, - .n_chunk_plans = count }; - status = coalesce_chunks(&dispatch, - (uint64_t)config->decode_workers * - config->max_encoded_chunk_bytes, - config->decode_workers, - indices, - scratch_reads); + 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); if (status != DAMACY_OK) goto Done; - reads.count = dispatch.n_read_ops; for (uint32_t i = 0; i < reads.count; ++i) reads.first_chunks[i] = UINT32_MAX; for (uint32_t i = count; i-- > 0;) { - uint32_t read = chunks[i].read_op_idx; + uint32_t read = read_index[i]; reads.chunks[i] = - (struct cpu_chunk_read){ .offset = chunks[i].offset_in_read, + (struct cpu_chunk_read){ .offset = offset_in_read[i], .next = reads.first_chunks[read] }; reads.first_chunks[read] = i; } *out = reads; Done: - free(chunks); free(scratch_reads); free(indices); if (status != DAMACY_OK) From 0dafafda1b1208b253db694b29cadf2745523efa Mon Sep 17 00:00:00 2001 From: Nathan Clack Date: Thu, 24 Sep 2026 21:13:48 +0000 Subject: [PATCH 3/3] test: fills beside a merged CPU read --- tests/test_cpu_executor.c | 36 ++++++++++++++++++++++++++++++++++++ 1 file changed, 36 insertions(+) diff --git a/tests/test_cpu_executor.c b/tests/test_cpu_executor.c index 75c4517c..c9c18255 100644 --- a/tests/test_cpu_executor.c +++ b/tests/test_cpu_executor.c @@ -389,6 +389,41 @@ test_single_worker_and_fills(void) return 0; } +static int +test_fills_share_merged_input(void) +{ + const struct test_chunk chunks[] = { + { .shard = 1, .offset = 16 }, { .missing = 1 }, + { .shard = 1, .offset = 24 }, { .missing = 1 }, + { .shard = 0, .offset = 0 }, { .shard = 0, .offset = 8 }, + }; + 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, 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); + struct damacy_batch* batch = NULL; + EXPECT(finish_batch(executor, &batch) == 0); + EXPECT(check_output(batch, &store, chunks, 6) == 0); + EXPECT(store.n_events == 1 && store.events[0].count == 2); + EXPECT(store.records[0].shard == 0 && store.records[0].offset == 0); + EXPECT(store.records[1].shard == 1 && store.records[1].offset == 16); + EXPECT(store.records[0].bytes == 2 * CHUNK_BYTES); + EXPECT(store.records[1].bytes == 2 * CHUNK_BYTES); + EXPECT(stats.waves_emitted == 1 && stats.reads_issued == 2); + EXPECT(stats.chunks_dispatched == 6 && stats.decode.count == 6); + EXPECT(stats.decode.input_bytes == 4 * CHUNK_BYTES); + EXPECT(stats.io.input_bytes == 4 * CHUNK_BYTES); + damacy_batch_release(batch); + damacy_executor_destroy(executor); + return 0; +} + static int test_batch_order_and_retained_output(void) { @@ -570,6 +605,7 @@ main(void) RUN(test_reader_limit_and_retry); RUN(test_overlapping_ranges); RUN(test_single_worker_and_fills); + RUN(test_fills_share_merged_input); RUN(test_batch_order_and_retained_output); RUN(test_plan_memory_budget); RUN(test_plan_memory_retry);