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
12 changes: 12 additions & 0 deletions bench/main.c
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down
34 changes: 34 additions & 0 deletions bench/scenarios/array-starve.json
Original file line number Diff line number Diff line change
@@ -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
}
}
18 changes: 18 additions & 0 deletions bench/sweep.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand All @@ -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}",
Expand All @@ -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):,}",
Expand Down
14 changes: 14 additions & 0 deletions src/damacy.h
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
13 changes: 10 additions & 3 deletions src/render_job/render_job.c
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
13 changes: 13 additions & 0 deletions src/render_job/render_job.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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
Expand Down
73 changes: 73 additions & 0 deletions src/wave/wave_input.c
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,78 @@
#include "log/log.h"
#include "wave_pool.h"

#include <pthread.h>
#include <stdio.h>
#include <stdlib.h>

// DAMACY_TRACE_WAVES=<file>: append one line per reserved wave:
// "<batch_id> <render_job_idx> <n_reads> <distinct_shards> <n_chunks>
// <input_bytes> <stop_reason>". 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)
{
Expand Down Expand Up @@ -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;
Expand Down
Loading