diff --git a/bench/main.c b/bench/main.c index c1aeba08..6750e8b4 100644 --- a/bench/main.c +++ b/bench/main.c @@ -929,6 +929,18 @@ emit_results(const struct scenario* sc, const struct run_metrics* rm, FILE* out) jw_uint(&jw, rm->stats.chunks_to_load); jw_key(&jw, "reads_issued"); jw_uint(&jw, rm->stats.reads_issued); + jw_key(&jw, "wave_reads_sum"); + jw_uint(&jw, rm->stats.wave_reads_sum); + jw_key(&jw, "wave_distinct_shards_sum"); + jw_uint(&jw, rm->stats.wave_distinct_shards_sum); + jw_key(&jw, "wave_stop_drained"); + jw_uint(&jw, rm->stats.wave_stop_drained); + jw_key(&jw, "wave_stop_host"); + jw_uint(&jw, rm->stats.wave_stop_host); + jw_key(&jw, "wave_stop_chunks"); + jw_uint(&jw, rm->stats.wave_stop_chunks); + jw_key(&jw, "wave_stop_dev"); + jw_uint(&jw, rm->stats.wave_stop_dev); jw_key(&jw, "distinct_zarrs"); jw_uint(&jw, rm->stats.array_meta.misses); jw_key(&jw, "distinct_shards"); diff --git a/bench/scenarios/array-starve.json b/bench/scenarios/array-starve.json new file mode 100644 index 00000000..b899a611 --- /dev/null +++ b/bench/scenarios/array-starve.json @@ -0,0 +1,34 @@ +{ + "name": "array-starve", + "dataset": { + "store_root": "/mnt/main0/home/nclack/data/damacy/array-starve", + "n_zarrs": 32, + "uri_fmt": "z%03u/scale0/image", + "zarr_shape": [256, 2048, 2048], + "chunk_shape": [16, 256, 256], + "shard_shape": [128, 2048, 2048], + "dtypes": ["u16"], + "codecs": ["zstd"], + "clevel": 3, + "entropy": 0.5, + "seed": 42 + }, + "sampling": { + "sample_shape": [64, 256, 256], + "n_batches": 60, + "n_warmup_batches": 10, + "samples_per_batch": 128, + "seed": 1234 + }, + "pipeline": { + "dtype": "f32", + "lookahead_samples": 512, + "n_io_threads": 64, + "metadata_io_concurrency": 64, + "max_gpu_memory_mb": 24576, + "n_array_meta_cache": 8192, + "n_shard_index_cache": 32768, + "n_chunk_layout_cache": 8192, + "max_shards_per_sample": 16 + } +} diff --git a/bench/sweep.py b/bench/sweep.py index bb559c11..430cd33c 100755 --- a/bench/sweep.py +++ b/bench/sweep.py @@ -111,6 +111,9 @@ def main( t.add_column("ttfb s", justify="right") t.add_column("wall s", justify="right") t.add_column("throughput GB/s", justify="right") + t.add_column("reads/wave", justify="right") + t.add_column("shards/wave", justify="right") + t.add_column("stop d/h/c/v", justify="right") t.add_column("lat max active", justify="right") t.add_column("lat sleep s", justify="right") t.add_column("meta read jobs", justify="right") @@ -130,6 +133,18 @@ def main( if input_transfer["ms_total"] > 0 else 0.0 ) + waves = c.get("waves_emitted", 0) or 1 + reads_per_wave = c.get("wave_reads_sum", 0) / waves + shards_per_wave = c.get("wave_distinct_shards_sum", 0) / waves + stops = "/".join( + str(c.get(k, 0)) + for k in ( + "wave_stop_drained", + "wave_stop_host", + "wave_stop_chunks", + "wave_stop_dev", + ) + ) t.add_row( str(v), f"{gb_in:.2f}", @@ -140,6 +155,9 @@ def main( f"{tm['time_to_first_batch'] / 1000.0:.2f}", f"{wall_ms / 1000.0:.2f}", f"{d['derived']['throughput_mb_s'] / 1e3:.2f}", + f"{reads_per_wave:.1f}", + f"{shards_per_wave:.1f}", + stops, f"{c.get('metadata_latency_max_active', 0):,}", f"{c.get('metadata_latency_total_sleep_ns', 0) / 1e9:.2f}", f"{c.get('metadata_backend_read_jobs', 0):,}", diff --git a/src/damacy.h b/src/damacy.h index cface9fe..4ca3267c 100644 --- a/src/damacy.h +++ b/src/damacy.h @@ -363,6 +363,20 @@ extern "C" uint64_t reads_issued; // real read_ops after coalesce uint64_t worker_steps; // scheduler ticks (proxy for worker CPU) + // Per-wave dispatch shape (probe for #154: does in-flight read + // concurrency track array/shard count instead of the io pool size?). + // Means are the sum divided by waves_emitted. distinct_shards_sum + // counts unique shard files within each wave's reads — the width the + // coalescer's round-robin spreads io across; a starved few-array run + // shows this far below n_io_threads. stop_* count why each wave + // stopped packing (see enum wave_stop_reason). + uint64_t wave_reads_sum; + uint64_t wave_distinct_shards_sum; + uint64_t wave_stop_drained; + uint64_t wave_stop_host; + uint64_t wave_stop_chunks; + uint64_t wave_stop_dev; + // Total GPU bytes currently committed to wave-resident buffers and // batch-output pools, counted against max_gpu_memory_bytes. Grows from // wave-init to the first damacy_pop (lazy batch pool sizing) and stays diff --git a/src/render_job/render_job.c b/src/render_job/render_job.c index 6f7888b8..9e7d1df0 100644 --- a/src/render_job/render_job.c +++ b/src/render_job/render_job.c @@ -203,17 +203,24 @@ wave_dispatcher_reserve(struct render_job* job, struct read_op_group_iterator it; read_op_group_iterator_init( &it, job->read_op_groups, job->n_read_op_groups, job->n_groups_dispatched); + out->stop_reason = WAVE_STOP_DRAINED; struct read_op_group g; while (read_op_group_iterator_next(&it, &g)) { struct read_op* r = &job->read_ops[g.read_op_idx]; int is_fill_group = job->chunk_plans[g.first_chunk].is_fill; uint64_t host_add = is_fill_group ? 0 : r->nbytes; - if (host_cursor + host_add > limits->input_cap) + if (host_cursor + host_add > limits->input_cap) { + out->stop_reason = WAVE_STOP_HOST; break; - if (take + g.n_chunks > limits->max_chunks_per_wave) + } + if (take + g.n_chunks > limits->max_chunks_per_wave) { + out->stop_reason = WAVE_STOP_CHUNKS; break; - if (dev_cursor + g.total_decompressed > limits->dev_decompressed_cap) + } + if (dev_cursor + g.total_decompressed > limits->dev_decompressed_cap) { + out->stop_reason = WAVE_STOP_DEV; break; + } uint64_t reserved_host_off = host_cursor; if (!is_fill_group) { diff --git a/src/render_job/render_job.h b/src/render_job/render_job.h index 280ff45f..b222759e 100644 --- a/src/render_job/render_job.h +++ b/src/render_job/render_job.h @@ -51,6 +51,18 @@ struct render_job_pool struct store_read; +// Why the wave packer stopped adding groups. Lets the bench tell a wave +// capped by a budget (more work was waiting) from one that drained the +// batch (no more work). Probe for #154: a starved few-array run shows +// mostly WAVE_STOP_DRAINED with few reads, not budget caps. +enum wave_stop_reason +{ + WAVE_STOP_DRAINED = 0, // consumed every remaining group + WAVE_STOP_HOST, // hit input_cap (host pinned buffer) + WAVE_STOP_CHUNKS, // hit max_chunks_per_wave + WAVE_STOP_DEV, // hit dev_decompressed_cap +}; + struct wave_desc { uint16_t render_job_idx; @@ -62,6 +74,7 @@ struct wave_desc uint64_t input_used_bytes; uint64_t io_bytes; uint8_t is_fill_wave; + uint8_t stop_reason; // enum wave_stop_reason }; struct wave_pack_limits diff --git a/src/wave/wave_input.c b/src/wave/wave_input.c index 0bcd8860..78348096 100644 --- a/src/wave/wave_input.c +++ b/src/wave/wave_input.c @@ -3,6 +3,78 @@ #include "log/log.h" #include "wave_pool.h" +#include +#include +#include + +// DAMACY_TRACE_WAVES=: append one line per reserved wave: +// " +// ". Probe for #154 — distinct_shards is the +// io width the coalescer's round-robin achieves for this wave. +static FILE* g_wave_trace; +static pthread_mutex_t g_wave_trace_lock = PTHREAD_MUTEX_INITIALIZER; +static pthread_once_t g_wave_trace_once = PTHREAD_ONCE_INIT; + +static void +wave_trace_open(void) +{ + const char* path = getenv("DAMACY_TRACE_WAVES"); + if (path && path[0]) + g_wave_trace = fopen(path, "a"); +} + +// Unique shard files among a wave's reads. The reads carry interned +// shard_path pointers (coalesce interns + round-robins them), so pointer +// identity is shard identity. O(n^2), but n_reads is small — and smallest +// exactly in the starved case this probe is chasing. +static uint32_t +wave_distinct_shards(const struct store_read* reads, uint32_t n) +{ + uint32_t distinct = 0; + for (uint32_t i = 0; i < n; ++i) { + int seen = 0; + for (uint32_t j = 0; j < i; ++j) + if (reads[j].key == reads[i].key) { + seen = 1; + break; + } + distinct += !seen; + } + return distinct; +} + +static void +wave_account(struct wave_pool* wp, + const struct input_slot* slot, + uint64_t batch_id, + const struct wave_desc* desc) +{ + uint32_t distinct = wave_distinct_shards(slot->store_reads, desc->n_reads); + wp->stats->wave_reads_sum += desc->n_reads; + wp->stats->wave_distinct_shards_sum += distinct; + switch (desc->stop_reason) { + case WAVE_STOP_HOST: wp->stats->wave_stop_host++; break; + case WAVE_STOP_CHUNKS: wp->stats->wave_stop_chunks++; break; + case WAVE_STOP_DEV: wp->stats->wave_stop_dev++; break; + default: wp->stats->wave_stop_drained++; break; + } + + pthread_once(&g_wave_trace_once, wave_trace_open); + if (g_wave_trace) { + pthread_mutex_lock(&g_wave_trace_lock); + fprintf(g_wave_trace, + "%llu %u %u %u %u %llu %u\n", + (unsigned long long)batch_id, + (unsigned)desc->render_job_idx, + (unsigned)desc->n_reads, + (unsigned)distinct, + (unsigned)desc->n_chunks, + (unsigned long long)desc->input_used_bytes, + (unsigned)desc->stop_reason); + pthread_mutex_unlock(&g_wave_trace_lock); + } +} + static void mark_changed(int* changed) { @@ -93,6 +165,7 @@ wave_input_reserve(struct wave_pool* wp, input_slot_begin_reservation(slot, &desc); wp->stats->waves_emitted++; wp->stats->chunks_dispatched += desc.n_chunks; + wave_account(wp, slot, job->batch_id, &desc); out->has_slot = 1; out->input_slot_idx = input_slot_idx;