diff --git a/src/frame.c b/src/frame.c index 9f33894..19a9c0e 100644 --- a/src/frame.c +++ b/src/frame.c @@ -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, diff --git a/src/hip_psf.h b/src/hip_psf.h index 3534a57..3700683 100644 --- a/src/hip_psf.h +++ b/src/hip_psf.h @@ -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); diff --git a/src/hip_psf.hip b/src/hip_psf.hip index dbefc9b..aa43c72 100644 --- a/src/hip_psf.hip +++ b/src/hip_psf.hip @@ -7,8 +7,50 @@ #include #include #include +#include +#include +#include +#include -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 host_refs; + std::vector host_tasks; + std::vector 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 counts, offsets, cursors, touched, host_refs; + std::vector host_tasks; + std::vector host_merges; + std::vector 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::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( + 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( + 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(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);