Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 14 additions & 7 deletions dev/cpu-pipeline.md
Original file line number Diff line number Diff line change
Expand Up @@ -74,20 +74,27 @@ 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.

## 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
Expand Down
22 changes: 17 additions & 5 deletions docs/pipeline.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. |
Expand All @@ -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
Expand Down
2 changes: 1 addition & 1 deletion src/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
114 changes: 77 additions & 37 deletions src/executor/coalesce.c
Original file line number Diff line number Diff line change
Expand Up @@ -2,41 +2,47 @@

#include "executor/read_op_sort.h"

#include <string.h>

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,
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 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];
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).
Expand All @@ -46,12 +52,12 @@ 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];
if (curr->shard_path == leader->shard_path &&
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 &&
leader_chunks < max_chunks_per_wave) {
Expand All @@ -62,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;
Expand All @@ -86,12 +92,12 @@ 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;
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) {
Expand All @@ -105,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;
Expand Down
23 changes: 19 additions & 4 deletions src/executor/coalesce.h
Original file line number Diff line number Diff line change
@@ -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.
Expand All @@ -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
Expand Down
Loading
Loading