Feat: Add producer-private adaptive HIP PSF prebinning

Move adaptive tile16 selection and support-expanded CPU binning into a
per-producer HipPsfPreparedChunk computed before acquiring the shared
sink lock. The sink keeps sole ownership of the GPU HDR, stream, device
scratch and upload staging; submit_prepared runs under the lock and
preserves the existing single-stream ordering and direct fallback path.
The legacy hip_psf_sink_submit entry point prepares an internal chunk so
single-threaded replay stays compatible, and default accumulation remains
atomic.
This commit is contained in:
wyj committed 2026-09-14 19:15:51 -04:00
1 parent 2d01c2c101
commit eec5dd215a
3 files changed
+531 -6

No files matched your search

+46 -3
View File
@@ -136,6 +136,8 @@ typedef struct {
const PsfKernelCache *cache;
#ifdef PSF_BACKEND_HIP
HipPsfSink *hip;
HipPsfPreparedChunk *prepared_chunk; /* Worker-owned CPU bin. */
int prepared_ready;
omp_lock_t *hip_lock; /* Shared sink lock; acquired only per chunk/fallback. */
int borrowed_hip;
double submit_seconds;
@@ -165,7 +167,8 @@ static int psf_event_sink_init(PsfEventSink *sink, double *hdr, int width,
if (cache != NULL) {
#ifdef PSF_BACKEND_DUMMY
if (dummy_psf_sink_create(&sink->dummy, width, height,
PSF_EVENT_SINK_CAPACITY)) {
PSF_EVENT_SINK_CAPACITY,
cache->max_radius_pixels)) {
fputs("Dummy PSF backend initialization failed\n", stderr);
dummy_psf_sink_destroy(sink->dummy);
sink->dummy = NULL;
@@ -198,6 +201,22 @@ static int psf_event_sink_init(PsfEventSink *sink, double *hdr, int width,
return 0;
}
#ifdef PSF_BACKEND_HIP
static int psf_event_sink_prepare(PsfEventSink *sink) {
if (sink->prepared_chunk == NULL || sink->count == 0 ||
sink->prepared_ready || sink->failed)
return sink->failed ? -1 : 0;
if (hip_psf_prepared_chunk_prepare(sink->prepared_chunk, sink->events,
sink->count, sink->hip_message,
sizeof sink->hip_message)) {
psf_event_sink_mark_failed(sink, "parallel event preparation");
return -1;
}
sink->prepared_ready = 1;
return 0;
}
#endif
static int psf_event_sink_flush_unlocked(PsfEventSink *sink) {
if (sink->failed)
return -1;
@@ -207,12 +226,18 @@ static int psf_event_sink_flush_unlocked(PsfEventSink *sink) {
#endif
#ifdef PSF_BACKEND_HIP
if (sink->hip != NULL) {
if (hip_psf_sink_submit(sink->hip, sink->events, sink->count,
sink->hip_message, sizeof sink->hip_message)) {
const int submit_result = sink->prepared_ready
? hip_psf_sink_submit_prepared(sink->hip, sink->prepared_chunk,
sink->events, sink->count,
sink->hip_message, sizeof sink->hip_message)
: hip_psf_sink_submit(sink->hip, sink->events, sink->count,
sink->hip_message, sizeof sink->hip_message);
if (submit_result) {
psf_event_sink_mark_failed(sink, "event submission");
return -1;
}
sink->count = 0;
sink->prepared_ready = 0;
return 0;
}
#endif
@@ -226,6 +251,7 @@ static int psf_event_sink_flush_unlocked(PsfEventSink *sink) {
static int psf_event_sink_flush(PsfEventSink *sink) {
#ifdef PSF_BACKEND_HIP
const double start = omp_get_wtime();
if (psf_event_sink_prepare(sink)) return -1;
if (sink->hip_lock) omp_set_lock(sink->hip_lock);
#endif
const int result = psf_event_sink_flush_unlocked(sink);
@@ -289,6 +315,7 @@ static int psf_event_sink_destroy(PsfEventSink *sink) {
#ifdef PSF_BACKEND_HIP
if (sink->borrowed_hip) {
const int result = psf_event_sink_flush(sink);
hip_psf_prepared_chunk_destroy(sink->prepared_chunk);
free(sink->events);
return result;
}
@@ -1184,6 +1211,7 @@ static int splat_catalog_tile(const Star *stars, size_t count,
#else
#ifdef PSF_BACKEND_HIP
const double submit_start = omp_get_wtime();
if (psf_event_sink_prepare(context->event_sink)) return -1;
if (context->event_sink->hip_lock)
omp_set_lock(context->event_sink->hip_lock);
#endif
@@ -1487,6 +1515,11 @@ size_t frame_splat_catalog(const FrameLensMesh *mesh,
snprintf(local.hip_message, sizeof local.hip_message, "worker event allocation failed");
psf_event_sink_mark_failed(&local, "producer initialization");
}
if (!local.failed && hip_psf_sink_is_adaptive(owner.hip) &&
hip_psf_prepared_chunk_create(&local.prepared_chunk, owner.hip,
local.hip_message,
sizeof local.hip_message))
psf_event_sink_mark_failed(&local, "producer preparation allocation");
const size_t worker = (size_t)omp_get_thread_num();
const size_t workers = (size_t)omp_get_num_threads();
size_t completed = 0, next_report = 8;
@@ -1524,6 +1557,16 @@ size_t frame_splat_catalog(const FrameLensMesh *mesh,
if (psf_event_sink_destroy(&owner)) hip_failed = 1;
fprintf(stderr, "HIP producers: %d workers; summed generation %.3f s, submission/fallback %.3f s; wall %.3f s\n",
hip_workers, hip_generate_seconds, hip_submit_seconds, omp_get_wtime() - hip_start);
fprintf(stderr,
"HIP PSF accumulation: atomic %zu batches, tile16 %zu batches; selector %.3f s, binning %.3f s\n",
owner.hip_timing.atomic_batch_count, owner.hip_timing.tile_batch_count,
owner.hip_timing.selection_seconds, owner.hip_timing.bin_seconds);
if (owner.hip_timing.tile_batch_count)
fprintf(stderr,
"HIP PSF tile workload: references %zu, tasks %zu, merges %zu; peak references/chunk %zu\n",
owner.hip_timing.tile_reference_count, owner.hip_timing.tile_task_count,
owner.hip_timing.tile_merge_count,
owner.hip_timing.peak_tile_reference_count);
copy_psf_splat_stats(psf_stats, (CatalogSplatStats){
.images = hip_images, .direct_fallbacks = hip_direct,
.cached_wing_clipped = hip_clipped, .discarded_below_min_y = hip_discarded,
+24
View File
@@ -12,10 +12,15 @@ extern "C" {
/* Opaque HIP state for one HDR framebuffer. It owns reusable device buffers
* for PsfCachedEvent chunks, immutable cache weights, and double RGB HDR. */
typedef struct HipPsfSink HipPsfSink;
typedef struct HipPsfPreparedChunk HipPsfPreparedChunk;
typedef struct {
size_t event_count, batch_count, timed_batch_count;
double upload_seconds, kernel_seconds, download_seconds;
size_t atomic_batch_count, tile_batch_count;
size_t tile_reference_count, tile_task_count, tile_merge_count;
size_t peak_tile_reference_count;
double selection_seconds, bin_seconds;
} HipPsfTiming;
int hip_psf_available(char *message, size_t message_size);
@@ -39,6 +44,25 @@ int hip_psf_sink_load_hdr(HipPsfSink *sink, const double *hdr,
int hip_psf_sink_submit(HipPsfSink *sink, const PsfCachedEvent *events,
size_t event_count, char *message, size_t message_size);
/* Optional adaptive path: each producer owns one preparation buffer. Prepare
* runs without the shared sink lock; submit_prepared runs under that lock and
* preserves the existing single-stream HDR ordering. Atomic sinks may keep
* using hip_psf_sink_submit directly. */
int hip_psf_sink_is_adaptive(const HipPsfSink *sink);
int hip_psf_prepared_chunk_create(HipPsfPreparedChunk **out,
const HipPsfSink *sink, char *message,
size_t message_size);
int hip_psf_prepared_chunk_prepare(HipPsfPreparedChunk *chunk,
const PsfCachedEvent *events,
size_t event_count, char *message,
size_t message_size);
int hip_psf_sink_submit_prepared(HipPsfSink *sink,
HipPsfPreparedChunk *chunk,
const PsfCachedEvent *events,
size_t event_count, char *message,
size_t message_size);
void hip_psf_prepared_chunk_destroy(HipPsfPreparedChunk *chunk);
/* Completes all work and overwrites hdr with double linear RGB device HDR. */
int hip_psf_sink_finish(HipPsfSink *sink, double *hdr, char *message,
size_t message_size);
+461 -3
View File
@@ -7,8 +7,50 @@
#include <limits>
#include <chrono>
#include <cstring>
#include <vector>
#include <algorithm>
#include <climits>
#include <cstdlib>
enum { PSF_LANES_PER_EVENT = 32, PSF_THREADS_PER_BLOCK = 128 };
enum { PSF_LANES_PER_EVENT = 32, PSF_THREADS_PER_BLOCK = 128,
TILE_SIZE = 16, TILE_TASK_EVENTS = 256 };
struct ProductionTileTask { unsigned tile, first, count, partial; };
struct ProductionTileMerge { unsigned tile, first, count; };
struct PreparedEvent {
int bx, by, support;
size_t phase_base;
double tx, ty, r, g, b;
};
struct HipTileScratch {
PsfCachedEvent *host = nullptr, *events = nullptr;
PreparedEvent *prepared = nullptr;
int *bounds = nullptr;
unsigned *refs = nullptr;
ProductionTileTask *tasks = nullptr;
ProductionTileMerge *merges = nullptr;
double *partials = nullptr;
hipEvent_t upload_start = nullptr, upload_end = nullptr;
hipEvent_t start = nullptr, end = nullptr;
size_t count = 0;
/* Submission-owned staging remains stable until collect_tile completes. */
std::vector<unsigned> host_refs;
std::vector<ProductionTileTask> host_tasks;
std::vector<ProductionTileMerge> host_merges;
};
struct HipPsfPreparedChunk {
const HipPsfSink *owner = nullptr;
size_t count = 0;
bool use_tile = false;
double selection_seconds = 0.0, bin_seconds = 0.0;
std::vector<unsigned> counts, offsets, cursors, touched, host_refs;
std::vector<ProductionTileTask> host_tasks;
std::vector<ProductionTileMerge> host_merges;
std::vector<unsigned> selector_marks;
unsigned selector_generation = 0;
};
/* Two reusable staging slots bound host/device queue memory. A slot may be
* overwritten only after its kernel_end event completes. All API calls for a
@@ -32,6 +74,9 @@ struct HipPsfSink {
double max_radius_pixels = 0.0;
hipEvent_t download_start = nullptr, download_end = nullptr;
HipPsfTiming timing = {};
bool adaptive = false;
HipTileScratch tile;
HipPsfPreparedChunk *legacy_prepared = nullptr;
std::chrono::steady_clock::time_point last_report = std::chrono::steady_clock::now();
};
@@ -59,6 +104,22 @@ static int collect_batch(HipPsfSink *sink, HipPsfBatch *batch,
return 0;
}
static int collect_tile(HipPsfSink *sink, char *message, size_t message_size) {
if (!sink->tile.count) return 0;
float upload_ms = 0.0f, kernel_ms = 0.0f;
if (report(hipEventSynchronize(sink->tile.end), message, message_size) ||
report(hipEventElapsedTime(&upload_ms, sink->tile.upload_start,
sink->tile.upload_end), message, message_size) ||
report(hipEventElapsedTime(&kernel_ms, sink->tile.start, sink->tile.end),
message, message_size)) return -1;
sink->timing.upload_seconds += upload_ms * 1e-3;
sink->timing.kernel_seconds += kernel_ms * 1e-3;
++sink->timing.timed_batch_count;
sink->completed_events += sink->tile.count;
sink->tile.count = 0;
return 0;
}
__device__ static size_t weight_index(int phase_resolution, int radius_pixels,
int phase_x, int phase_y, int offset_x, int offset_y) {
const size_t nodes = (size_t)phase_resolution + 1;
@@ -121,6 +182,94 @@ __global__ static void splat_kernel(const PsfCachedEvent *events, size_t count,
}
}
__global__ static void prepare_tile_events(
const PsfCachedEvent *events, size_t count, int phases, int radius,
double max_radius, PreparedEvent *prepared, int *bounds) {
const size_t thread = (size_t)blockIdx.x * blockDim.x + threadIdx.x;
const size_t index = thread / PSF_LANES_PER_EVENT;
const int lane = thread % PSF_LANES_PER_EVENT;
if (index >= count) return;
const PsfCachedEvent event = events[index];
const int bx = (int)floor(event.x), by = (int)floor(event.y);
const double fx = event.x - bx, fy = event.y - by;
const double ax = fx * phases, ay = fy * phases;
const int x0 = (int)floor(ax), y0 = (int)floor(ay);
const int support = (int)fmin(ceil(event.support_radius), max_radius);
if (!lane)
prepared[index] = {bx, by, support,
weight_index(phases, radius, x0, y0, 0, 0), ax - x0, ay - y0,
event.color.r * event.flux, event.color.g * event.flux,
event.color.b * event.flux};
const int side = 2 * radius + 1;
for (int row = lane; row < side; row += PSF_LANES_PER_EVENT) {
int first = 1, last = 0;
const int dy = row - radius;
if (dy >= -support && dy <= support)
row_range(event.support_radius, fx, fy, support, dy, &first, &last);
bounds[2 * (index * side + row)] = first;
bounds[2 * (index * side + row) + 1] = last;
}
}
__global__ static void accumulate_tiles(
const PreparedEvent *events, const int *bounds, const unsigned *refs,
const ProductionTileTask *tasks, int tiles_x, int width, int height,
const float *weights, int phases, int radius, double *partials,
double *hdr) {
const ProductionTileTask task = tasks[blockIdx.x];
const int local = threadIdx.x;
const int px = (task.tile % tiles_x) * TILE_SIZE + local % TILE_SIZE;
const int py = (task.tile / tiles_x) * TILE_SIZE + local / TILE_SIZE;
const int side = 2 * radius + 1;
const size_t plane = (size_t)side * side;
double r = 0.0, g = 0.0, b = 0.0;
if (px < width && py < height) for (unsigned k = 0; k < task.count; ++k) {
const unsigned id = refs[task.first + k];
const PreparedEvent event = events[id];
const int dx = px - event.bx, dy = py - event.by;
if (dy < -event.support || dy > event.support) continue;
const size_t row = 2 * ((size_t)id * side + dy + radius);
if (dx < bounds[row] || dx > bounds[row + 1]) continue;
const size_t p = event.phase_base + (ptrdiff_t)dy * side + dx;
const double w00 = weights[p], w10 = weights[p + plane];
const double w01 = weights[p + (phases + 1) * plane];
const double w11 = weights[p + (phases + 2) * plane];
const double weight = (1.0 - event.ty) *
((1.0 - event.tx) * w00 + event.tx * w10) +
event.ty *
((1.0 - event.tx) * w01 + event.tx * w11);
r += event.r * weight; g += event.g * weight; b += event.b * weight;
}
if (task.partial == UINT_MAX) {
if (px < width && py < height) {
const size_t p = 3 * ((size_t)py * width + px);
hdr[p] += r; hdr[p + 1] += g; hdr[p + 2] += b;
}
} else {
const size_t p = 3 * ((size_t)task.partial * TILE_SIZE * TILE_SIZE + local);
partials[p] = r; partials[p + 1] = g; partials[p + 2] = b;
}
}
__global__ static void merge_tiles(const ProductionTileMerge *merges,
const ProductionTileTask *tasks, int tiles_x,
int width, int height,
const double *partials, double *hdr) {
const ProductionTileMerge merge = merges[blockIdx.x];
const int local = threadIdx.x;
const int px = (merge.tile % tiles_x) * TILE_SIZE + local % TILE_SIZE;
const int py = (merge.tile / tiles_x) * TILE_SIZE + local / TILE_SIZE;
if (px >= width || py >= height) return;
double r = 0.0, g = 0.0, b = 0.0;
for (unsigned k = 0; k < merge.count; ++k) {
const ProductionTileTask task = tasks[merge.first + k];
const size_t p = 3 * ((size_t)task.partial * TILE_SIZE * TILE_SIZE + local);
r += partials[p]; g += partials[p + 1]; b += partials[p + 2];
}
const size_t p = 3 * ((size_t)py * width + px);
hdr[p] += r; hdr[p + 1] += g; hdr[p + 2] += b;
}
extern "C" int hip_psf_available(char *message, size_t message_size) {
int count = 0;
if (report(hipGetDeviceCount(&count), message, message_size)) return -1;
@@ -157,6 +306,23 @@ extern "C" int hip_psf_sink_create(HipPsfSink **out, int width, int height,
sink->hdr_values = (size_t)width * (size_t)height * 3;
sink->width = width; sink->height = height; sink->phase_resolution = cache->phase_resolution;
sink->radius_pixels = cache->radius_pixels; sink->max_radius_pixels = cache->max_radius_pixels;
const char *accumulation = std::getenv("HIP_PSF_ACCUMULATION");
if (accumulation && std::strcmp(accumulation, "atomic") &&
std::strcmp(accumulation, "adaptive")) {
if (message && message_size)
std::snprintf(message, message_size,
"HIP_PSF_ACCUMULATION must be atomic or adaptive");
hip_psf_sink_destroy(sink); return -1;
}
sink->adaptive = accumulation && !std::strcmp(accumulation, "adaptive");
if (sink->adaptive &&
(event_capacity > UINT_MAX ||
event_capacity > std::numeric_limits<size_t>::max() / side / 2 /
sizeof(int))) {
if (message && message_size)
std::snprintf(message, message_size, "adaptive HIP PSF size overflow");
hip_psf_sink_destroy(sink); return -1;
}
const size_t weight_count = nodes * nodes * side * side;
if (report(hipGetDevice(&sink->device), message, message_size) ||
report(hipStreamCreate(&sink->stream), message, message_size) ||
@@ -178,6 +344,37 @@ extern "C" int hip_psf_sink_create(HipPsfSink **out, int width, int height,
hip_psf_sink_destroy(sink); return -1;
}
}
if (sink->adaptive) {
const size_t tiles_x = ((size_t)width + TILE_SIZE - 1) / TILE_SIZE;
const size_t tiles_y = ((size_t)height + TILE_SIZE - 1) / TILE_SIZE;
const size_t tile_count = tiles_x * tiles_y;
const size_t bounds_count = event_capacity * side * 2;
const size_t ref_bytes = 32u * 1024u * 1024u;
const size_t task_bytes = 2u * 1024u * 1024u;
const size_t partial_bytes = 128u * 1024u * 1024u;
if (hip_psf_prepared_chunk_create(&sink->legacy_prepared, sink,
message, message_size)) {
hip_psf_sink_destroy(sink); return -1;
}
if (report(hipHostMalloc(&sink->tile.host,
event_capacity * sizeof(PsfCachedEvent)), message, message_size) ||
report(hipMalloc(&sink->tile.events,
event_capacity * sizeof(PsfCachedEvent)), message, message_size) ||
report(hipMalloc(&sink->tile.prepared,
event_capacity * sizeof(PreparedEvent)), message, message_size) ||
report(hipMalloc(&sink->tile.bounds, bounds_count * sizeof(int)), message, message_size) ||
report(hipMalloc(&sink->tile.refs, ref_bytes), message, message_size) ||
report(hipMalloc(&sink->tile.tasks, task_bytes), message, message_size) ||
report(hipMalloc(&sink->tile.merges,
tile_count * sizeof(ProductionTileMerge)), message, message_size) ||
report(hipMalloc(&sink->tile.partials, partial_bytes), message, message_size) ||
report(hipEventCreate(&sink->tile.upload_start), message, message_size) ||
report(hipEventCreate(&sink->tile.upload_end), message, message_size) ||
report(hipEventCreate(&sink->tile.start), message, message_size) ||
report(hipEventCreate(&sink->tile.end), message, message_size)) {
hip_psf_sink_destroy(sink); return -1;
}
}
/* The caller may release the immutable host cache after creation. */
if (report(hipStreamSynchronize(sink->stream), message, message_size)) {
hip_psf_sink_destroy(sink); return -1;
@@ -185,14 +382,239 @@ extern "C" int hip_psf_sink_create(HipPsfSink **out, int width, int height,
*out = sink; ok(message, message_size); return 0;
}
extern "C" int hip_psf_sink_submit(HipPsfSink *sink, const PsfCachedEvent *events,
size_t event_count, char *message, size_t message_size) {
extern "C" int hip_psf_sink_is_adaptive(const HipPsfSink *sink) {
return sink && sink->adaptive;
}
extern "C" int hip_psf_prepared_chunk_create(HipPsfPreparedChunk **out,
const HipPsfSink *sink, char *message, size_t message_size) {
if (!out || !sink || !sink->adaptive) {
if (message && message_size) std::snprintf(message, message_size,
"invalid adaptive chunk creation");
return -1;
}
*out = nullptr;
try {
HipPsfPreparedChunk *chunk = new HipPsfPreparedChunk;
*out = chunk;
chunk->owner = sink;
const size_t tile_count = ((size_t)sink->width + TILE_SIZE - 1) / TILE_SIZE *
(((size_t)sink->height + TILE_SIZE - 1) / TILE_SIZE);
chunk->counts.resize(tile_count);
chunk->offsets.resize(tile_count);
chunk->cursors.resize(tile_count);
chunk->selector_marks.resize(((size_t)sink->width + 31) / 32 *
(((size_t)sink->height + 31) / 32));
*out = chunk;
} catch (...) {
delete *out;
*out = nullptr;
if (message && message_size) std::snprintf(message, message_size,
"adaptive chunk host allocation failed");
return -1;
}
ok(message, message_size); return 0;
}
extern "C" void hip_psf_prepared_chunk_destroy(HipPsfPreparedChunk *chunk) {
delete chunk;
}
extern "C" int hip_psf_prepared_chunk_prepare(HipPsfPreparedChunk *chunk,
const PsfCachedEvent *events, size_t event_count,
char *message, size_t message_size) {
if (!chunk || !chunk->owner || (event_count && !events) ||
event_count > chunk->owner->event_capacity) {
if (message && message_size) std::snprintf(message, message_size,
"invalid adaptive chunk");
return -1;
}
const HipPsfSink *sink = chunk->owner;
chunk->count = event_count;
chunk->use_tile = false;
chunk->selection_seconds = chunk->bin_seconds = 0.0;
if (event_count < 8192) { ok(message, message_size); return 0; }
const auto select_start = std::chrono::steady_clock::now();
if (++chunk->selector_generation == 0) {
std::fill(chunk->selector_marks.begin(), chunk->selector_marks.end(), 0u);
chunk->selector_generation = 1;
}
const int selector_x = (sink->width + 31) / 32;
size_t occupied = 0;
for (size_t i = 0; i < event_count; ++i) {
const int x = (int)floor(events[i].x), y = (int)floor(events[i].y);
if (x < 0 || x >= sink->width || y < 0 || y >= sink->height) continue;
unsigned &mark = chunk->selector_marks[(size_t)(y / 32) * selector_x + x / 32];
if (mark != chunk->selector_generation) {
mark = chunk->selector_generation;
++occupied;
}
}
bool use_tile = occupied && event_count / occupied >= 32;
chunk->selection_seconds = std::chrono::duration<double>(
std::chrono::steady_clock::now() - select_start).count();
if (!use_tile) { ok(message, message_size); return 0; }
const auto bin_start = std::chrono::steady_clock::now();
for (unsigned index : chunk->touched) chunk->counts[index] = 0;
chunk->touched.clear(); chunk->host_tasks.clear(); chunk->host_merges.clear();
const int tiles_x = (sink->width + TILE_SIZE - 1) / TILE_SIZE;
const size_t ref_capacity = 32u * 1024u * 1024u / sizeof(unsigned);
const size_t task_capacity = 2u * 1024u * 1024u / sizeof(ProductionTileTask);
const size_t partial_capacity = 128u * 1024u * 1024u /
(TILE_SIZE * TILE_SIZE * 3 * sizeof(double));
size_t refs_count = 0;
try {
for (size_t i = 0; i < event_count && use_tile; ++i) {
const int support = (int)fmin(ceil(events[i].support_radius), sink->max_radius_pixels);
const int bx = (int)floor(events[i].x), by = (int)floor(events[i].y);
const int x0 = std::max(0, bx - support), x1 = std::min(sink->width - 1, bx + support);
const int y0 = std::max(0, by - support), y1 = std::min(sink->height - 1, by + support);
if (x0 > x1 || y0 > y1) continue;
for (int y = y0 / TILE_SIZE; y <= y1 / TILE_SIZE && use_tile; ++y)
for (int x = x0 / TILE_SIZE; x <= x1 / TILE_SIZE && use_tile; ++x) {
const unsigned index = (unsigned)(y * tiles_x + x);
if (!chunk->counts[index]) chunk->touched.push_back(index);
++chunk->counts[index];
if (++refs_count > ref_capacity) { use_tile = false; break; }
}
}
size_t partial_count = 0;
if (use_tile) {
chunk->host_refs.resize(refs_count);
size_t offset = 0;
for (unsigned index : chunk->touched) {
chunk->offsets[index] = offset;
chunk->cursors[index] = offset;
offset += chunk->counts[index];
}
for (size_t i = 0; i < event_count; ++i) {
const int support = (int)fmin(ceil(events[i].support_radius), sink->max_radius_pixels);
const int bx = (int)floor(events[i].x), by = (int)floor(events[i].y);
const int x0 = std::max(0, bx - support), x1 = std::min(sink->width - 1, bx + support);
const int y0 = std::max(0, by - support), y1 = std::min(sink->height - 1, by + support);
if (x0 > x1 || y0 > y1) continue;
for (int y = y0 / TILE_SIZE; y <= y1 / TILE_SIZE; ++y)
for (int x = x0 / TILE_SIZE; x <= x1 / TILE_SIZE; ++x) {
const unsigned index = (unsigned)(y * tiles_x + x);
chunk->host_refs[chunk->cursors[index]++] = (unsigned)i;
}
}
for (unsigned index : chunk->touched) {
const size_t count = chunk->counts[index];
const unsigned first_task = (unsigned)chunk->host_tasks.size();
if (chunk->host_tasks.size() + (count + TILE_TASK_EVENTS - 1) / TILE_TASK_EVENTS >
task_capacity) { use_tile = false; break; }
for (size_t k = 0; k < count; k += TILE_TASK_EVENTS) {
const unsigned length = (unsigned)std::min((size_t)TILE_TASK_EVENTS, count - k);
const unsigned partial = count <= TILE_TASK_EVENTS ? UINT_MAX : (unsigned)partial_count++;
if (partial_count > partial_capacity) { use_tile = false; break; }
chunk->host_tasks.push_back({index, (unsigned)(chunk->offsets[index] + k), length, partial});
}
if (!use_tile) break;
if (count > TILE_TASK_EVENTS)
chunk->host_merges.push_back({index, first_task,
(unsigned)(chunk->host_tasks.size() - first_task)});
}
}
} catch (...) {
if (message && message_size) std::snprintf(message, message_size,
"adaptive chunk bin allocation failed");
return -1;
}
chunk->bin_seconds = std::chrono::duration<double>(
std::chrono::steady_clock::now() - bin_start).count();
chunk->use_tile = use_tile && !chunk->host_tasks.empty();
ok(message, message_size); return 0;
}
static int submit_prepared(HipPsfSink *sink, HipPsfPreparedChunk *prepared,
const PsfCachedEvent *events, size_t event_count,
char *message, size_t message_size) {
if (!sink || (event_count && !events) || event_count > sink->event_capacity) {
if (message && message_size) std::snprintf(message, message_size, "invalid HIP PSF event chunk");
return -1;
}
if (!event_count) { ok(message, message_size); return 0; }
if (report(hipSetDevice(sink->device), message, message_size)) return -1;
if (sink->adaptive &&
(!prepared || prepared->owner != sink || prepared->count != event_count)) {
if (message && message_size) std::snprintf(message, message_size,
"adaptive chunk not prepared for sink/count");
return -1;
}
if (sink->adaptive) {
sink->timing.selection_seconds += prepared->selection_seconds;
sink->timing.bin_seconds += prepared->bin_seconds;
}
const bool use_tile = sink->adaptive && prepared->use_tile;
if (use_tile) {
HipTileScratch &tile = sink->tile;
if (collect_tile(sink, message, message_size)) return -1;
const int tiles_x = (sink->width + TILE_SIZE - 1) / TILE_SIZE;
try {
tile.host_refs = prepared->host_refs;
tile.host_tasks = prepared->host_tasks;
tile.host_merges = prepared->host_merges;
} catch (...) {
if (message && message_size) std::snprintf(message, message_size,
"adaptive tile staging allocation failed");
return -1;
}
{
std::memcpy(tile.host, events, event_count * sizeof(PsfCachedEvent));
if (report(hipEventRecord(tile.upload_start, sink->stream), message, message_size) ||
report(hipMemcpyAsync(tile.events, tile.host,
event_count * sizeof(PsfCachedEvent), hipMemcpyHostToDevice,
sink->stream), message, message_size) ||
report(hipMemcpyAsync(tile.refs, tile.host_refs.data(),
tile.host_refs.size() * sizeof(unsigned), hipMemcpyHostToDevice,
sink->stream), message, message_size) ||
report(hipMemcpyAsync(tile.tasks, tile.host_tasks.data(),
tile.host_tasks.size() * sizeof(ProductionTileTask), hipMemcpyHostToDevice,
sink->stream), message, message_size) ||
(!tile.host_merges.empty() && report(hipMemcpyAsync(
tile.merges, tile.host_merges.data(),
tile.host_merges.size() * sizeof(ProductionTileMerge), hipMemcpyHostToDevice,
sink->stream), message, message_size)) ||
report(hipEventRecord(tile.upload_end, sink->stream), message, message_size) ||
report(hipEventRecord(tile.start, sink->stream), message, message_size)) return -1;
const size_t prepare_threads = event_count * PSF_LANES_PER_EVENT;
hipLaunchKernelGGL(prepare_tile_events,
dim3((unsigned)((prepare_threads + PSF_THREADS_PER_BLOCK - 1) /
PSF_THREADS_PER_BLOCK)), dim3(PSF_THREADS_PER_BLOCK), 0,
sink->stream, tile.events, event_count, sink->phase_resolution,
sink->radius_pixels, sink->max_radius_pixels, tile.prepared, tile.bounds);
hipLaunchKernelGGL(accumulate_tiles, dim3((unsigned)prepared->host_tasks.size()),
dim3(TILE_SIZE * TILE_SIZE), 0, sink->stream, tile.prepared, tile.bounds,
tile.refs, tile.tasks, tiles_x, sink->width, sink->height, sink->weights,
sink->phase_resolution, sink->radius_pixels, tile.partials, sink->hdr);
if (!prepared->host_merges.empty())
hipLaunchKernelGGL(merge_tiles, dim3((unsigned)prepared->host_merges.size()),
dim3(TILE_SIZE * TILE_SIZE), 0, sink->stream, tile.merges, tile.tasks,
tiles_x, sink->width, sink->height, tile.partials, sink->hdr);
if (report(hipGetLastError(), message, message_size) ||
report(hipEventRecord(tile.end, sink->stream), message, message_size)) return -1;
tile.count = event_count;
sink->timing.tile_reference_count += prepared->host_refs.size();
sink->timing.tile_task_count += prepared->host_tasks.size();
sink->timing.tile_merge_count += prepared->host_merges.size();
sink->timing.peak_tile_reference_count = std::max(
sink->timing.peak_tile_reference_count, prepared->host_refs.size());
++sink->timing.tile_batch_count;
++sink->timing.batch_count;
sink->timing.event_count += event_count;
const auto now = std::chrono::steady_clock::now();
if (std::chrono::duration<double>(now - sink->last_report).count() >= 5.0) {
std::fprintf(stderr,
"HIP PSF progress: %zu submitted / %zu completed events; %zu / %zu batches timed; kernel %.3f s\n",
sink->timing.event_count, sink->completed_events,
sink->timing.timed_batch_count, sink->timing.batch_count,
sink->timing.kernel_seconds);
sink->last_report = now;
}
ok(message, message_size); return 0;
}
}
const unsigned int threads = PSF_THREADS_PER_BLOCK;
const size_t per_block = PSF_THREADS_PER_BLOCK / PSF_LANES_PER_EVENT;
const size_t blocks = event_count / per_block + (event_count % per_block != 0);
@@ -214,6 +636,7 @@ extern "C" int hip_psf_sink_submit(HipPsfSink *sink, const PsfCachedEvent *event
if (report(hipGetLastError(), message, message_size) ||
report(hipEventRecord(batch.kernel_end, sink->stream), message, message_size)) return -1;
batch.count = event_count;
++sink->timing.atomic_batch_count;
++sink->timing.batch_count;
sink->timing.event_count += event_count;
const auto now = std::chrono::steady_clock::now();
@@ -226,6 +649,26 @@ extern "C" int hip_psf_sink_submit(HipPsfSink *sink, const PsfCachedEvent *event
ok(message, message_size); return 0;
}
extern "C" int hip_psf_sink_submit_prepared(HipPsfSink *sink,
HipPsfPreparedChunk *chunk, const PsfCachedEvent *events,
size_t event_count, char *message, size_t message_size) {
return submit_prepared(sink, chunk, events, event_count, message, message_size);
}
extern "C" int hip_psf_sink_submit(HipPsfSink *sink,
const PsfCachedEvent *events, size_t event_count,
char *message, size_t message_size) {
if (!sink) {
if (message && message_size) std::snprintf(message, message_size,
"invalid HIP PSF sink");
return -1;
}
if (sink->adaptive && hip_psf_prepared_chunk_prepare(
sink->legacy_prepared, events, event_count, message, message_size)) return -1;
return submit_prepared(sink, sink->legacy_prepared, events, event_count,
message, message_size);
}
extern "C" int hip_psf_sink_load_hdr(HipPsfSink *sink, const double *hdr,
char *message, size_t message_size) {
if (!sink || !hdr) {
@@ -233,6 +676,7 @@ extern "C" int hip_psf_sink_load_hdr(HipPsfSink *sink, const double *hdr,
return -1;
}
if (report(hipSetDevice(sink->device), message, message_size) ||
collect_tile(sink, message, message_size) ||
report(hipMemcpyAsync(sink->hdr, hdr, sink->hdr_values * sizeof *hdr,
hipMemcpyHostToDevice, sink->stream), message, message_size) ||
report(hipStreamSynchronize(sink->stream), message, message_size))
@@ -248,6 +692,7 @@ extern "C" int hip_psf_sink_finish(HipPsfSink *sink, double *hdr, char *message,
report(hipEventRecord(sink->download_end, sink->stream), message, message_size) ||
report(hipStreamSynchronize(sink->stream), message, message_size)) return -1;
float milliseconds = 0.0f;
if (collect_tile(sink, message, message_size)) return -1;
for (HipPsfBatch &batch : sink->slots)
if (collect_batch(sink, &batch, message, message_size)) return -1;
if (report(hipEventElapsedTime(&milliseconds, sink->download_start, sink->download_end),
@@ -264,6 +709,7 @@ extern "C" int hip_psf_sink_get_timing(const HipPsfSink *sink, HipPsfTiming *tim
extern "C" void hip_psf_sink_destroy(HipPsfSink *sink) {
if (!sink) return;
hip_psf_prepared_chunk_destroy(sink->legacy_prepared);
(void)hipSetDevice(sink->device);
if (sink->stream) (void)hipStreamSynchronize(sink->stream);
for (HipPsfBatch &batch : sink->slots) {
@@ -274,6 +720,18 @@ extern "C" void hip_psf_sink_destroy(HipPsfSink *sink) {
if (batch.host) (void)hipHostFree(batch.host);
if (batch.device) (void)hipFree(batch.device);
}
if (sink->tile.upload_start) (void)hipEventDestroy(sink->tile.upload_start);
if (sink->tile.upload_end) (void)hipEventDestroy(sink->tile.upload_end);
if (sink->tile.start) (void)hipEventDestroy(sink->tile.start);
if (sink->tile.end) (void)hipEventDestroy(sink->tile.end);
if (sink->tile.host) (void)hipHostFree(sink->tile.host);
if (sink->tile.events) (void)hipFree(sink->tile.events);
if (sink->tile.prepared) (void)hipFree(sink->tile.prepared);
if (sink->tile.bounds) (void)hipFree(sink->tile.bounds);
if (sink->tile.refs) (void)hipFree(sink->tile.refs);
if (sink->tile.tasks) (void)hipFree(sink->tile.tasks);
if (sink->tile.merges) (void)hipFree(sink->tile.merges);
if (sink->tile.partials) (void)hipFree(sink->tile.partials);
if (sink->download_start) (void)hipEventDestroy(sink->download_start);
if (sink->download_end) (void)hipEventDestroy(sink->download_end);
if (sink->weights) (void)hipFree(sink->weights);