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
13 changes: 13 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -20214,6 +20214,19 @@ used consistently for `s->adm_csf_module`, `s->adm_csf_den_module`, and
removing the dead code path entirely.


- CUDA: the ADR-0242 drain batch is scoped to the engine that opened it. It was
thread-local with no owner, so two `VmafContext`s on one OS thread shared it: a
closed context left its `CUevent`s and its `bool *` drained flags registered and
the next context waited on freed objects. Registration by a foreign owner is now
refused, a flush consumes its entries, and `vmaf_close()` clears them.
- `vmaf_score_at_index()` and `vmaf_feature_score_at_index()` no longer read the
feature collector without a fence (Netflix/vmaf#1305): when the lock-free read
reports the slot unwritten they wait for the worker threads, drain the CUDA batch
and run the pending collect for indices at or below the requested one, then read
again. Reads that hit written slots are unchanged, so the streaming path keeps
its throughput.


- **`--feature ssim --backend cuda` silently broken since introduction.**
`integer_ssim/integer_ssim_score.cu` defined three `__global__`
kernels (`integer_ssim_horiz_8bpc`, `integer_ssim_horiz_16bpc`,
Expand Down
11 changes: 11 additions & 0 deletions changelog.d/fixed/cuda-drain-batch-per-state-lifetime.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
- CUDA: the ADR-0242 drain batch is scoped to the engine that opened it. It was
thread-local with no owner, so two `VmafContext`s on one OS thread shared it: a
closed context left its `CUevent`s and its `bool *` drained flags registered and
the next context waited on freed objects. Registration by a foreign owner is now
refused, a flush consumes its entries, and `vmaf_close()` clears them.
- `vmaf_score_at_index()` and `vmaf_feature_score_at_index()` no longer read the
feature collector without a fence (Netflix/vmaf#1305): when the lock-free read
reports the slot unwritten they wait for the worker threads, drain the CUDA batch
and run the pending collect for indices at or below the requested one, then read
again. Reads that hit written slots are unchanged, so the streaming path keeps
its throughput.
12 changes: 12 additions & 0 deletions core/src/cuda/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -243,6 +243,18 @@ cuda/
[ADR-0356](../../../docs/adr/0356-ssimulacra2-cuda-leaks-perf.md)
(`ssimulacra2_cuda` had two unloaded modules — `module_blur` +
`module_mul`).
- **The drain batch belongs to one engine at a time.** `drain_batch.c`'s
`g_drain_batch` is thread-local (ADR-0242) but two `VmafContext`s can run on
one OS thread, so it carries the owning `VmafCudaState`:
`vmaf_cuda_drain_batch_open()` takes that state and drops entries left by a
different owner, `vmaf_cuda_drain_batch_flush()` returns 0 without touching
CUDA when the caller is not the owner, the flush clears the entries it
consumed, and `vmaf_cuda_drain_batch_thread_destroy()` wipes entries, the open
flag and the owner before the engine state is freed. Never restore the
owner-less `open(void)` signature: without it a closed context leaves dangling
`CUevent`s and freed `bool *` flags for the next one (docs/state.md
T-UPSTREAM-1305-CUDA-DRAIN-BATCH-THREAD-GLOBAL-2026-09-03, pinned by
`core/test/test_cuda_drain_batch.c`).

## Governing ADRs

Expand Down
40 changes: 39 additions & 1 deletion core/src/cuda/drain_batch.c
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,13 @@
* the ``thread_pool`` only parallelises CPU extractors), so a TLS
* batch matches the actual call graph 1:1. */
typedef struct DrainBatchTls {
/* Owning engine state. The batch is thread-local (ADR-0242) but two
* VmafContexts can share one OS thread: without an owner check, context
* B's flush would wait on context A's CUevents and write through A's
* ``bool *`` flags after A was closed. Registration is refused for a
* foreign owner, and open()/thread_destroy() re-claim / clear it.
* docs/state.md T-UPSTREAM-1305-CUDA-DRAIN-BATCH-THREAD-GLOBAL-2026-09-03. */
const VmafCudaState *owner;
bool open;
unsigned n;
CUevent finished[VMAF_CUDA_DRAIN_BATCH_MAX];
Expand All @@ -48,8 +55,17 @@ typedef struct DrainBatchTls {

static _Thread_local DrainBatchTls g_drain_batch;

void vmaf_cuda_drain_batch_open(void)
void vmaf_cuda_drain_batch_open(const VmafCudaState *cu_state)
{
if (g_drain_batch.owner != cu_state) {
/* A different engine (or the first one on this thread) is taking the
* batch over. Any entries left by the previous owner refer to events
* and flags whose lifetime we cannot vouch for, so drop them rather
* than wait on them. The previous owner's own close()/destroy() has
* already released the objects themselves. */
g_drain_batch.n = 0;
g_drain_batch.owner = cu_state;
}
if (g_drain_batch.open) {
return;
}
Expand Down Expand Up @@ -137,6 +153,10 @@ int vmaf_cuda_drain_batch_flush(VmafCudaState *cu_state)
if (cu_state == NULL) {
return -EINVAL;
}
if (g_drain_batch.owner != NULL && g_drain_batch.owner != cu_state) {
/* Another engine owns this thread's batch — never wait on its events. */
return 0;
}
if (!g_drain_batch.open || g_drain_batch.n == 0U) {
return 0;
}
Expand Down Expand Up @@ -173,6 +193,10 @@ int vmaf_cuda_drain_batch_flush(VmafCudaState *cu_state)
*g_drain_batch.flags[i] = true;
}
}
/* Entries are consumed: the flags are set and the events belong to the
* frame that just drained. Clearing here means a later flush without an
* intervening open() cannot wait on them a second time. */
g_drain_batch.n = 0;
return 0;
}

Expand All @@ -186,8 +210,22 @@ void vmaf_cuda_drain_batch_close(void)
* paths in integer_motion/adm/vif_cuda.c. */
}

unsigned vmaf_cuda_drain_batch_pending(void)
{
return g_drain_batch.n;
}

void vmaf_cuda_drain_batch_thread_destroy(VmafCudaState *cu_state)
{
/* Wipe the registrations before anything else: after this call the
* caller frees the engine state, so every registered CUevent and every
* ``bool *`` flag becomes dangling. A later context on the same thread
* must not see them (T-UPSTREAM-1305-CUDA-DRAIN-BATCH-THREAD-GLOBAL). */
if (g_drain_batch.owner == NULL || g_drain_batch.owner == cu_state) {
g_drain_batch.n = 0;
g_drain_batch.open = false;
g_drain_batch.owner = NULL;
}
if (cu_state == NULL || g_drain_batch.drain_str == NULL) {
g_drain_batch.drain_str = NULL;
return;
Expand Down
14 changes: 13 additions & 1 deletion core/src/cuda/drain_batch.h
Original file line number Diff line number Diff line change
Expand Up @@ -80,7 +80,7 @@ struct VmafCudaKernelLifecycle;
* Idempotent: a second open without an intervening close is a no-op
* and keeps the existing batch.
*/
void vmaf_cuda_drain_batch_open(void);
void vmaf_cuda_drain_batch_open(const VmafCudaState *cu_state);

/**
* Register an extractor lifecycle into the open batch.
Expand Down Expand Up @@ -150,6 +150,18 @@ void vmaf_cuda_drain_batch_close(void);
*/
void vmaf_cuda_drain_batch_thread_destroy(VmafCudaState *cu_state);

/**
* Number of entries currently registered in this thread's batch.
*
* Introspection for the unit test that pins the owner scoping
* (docs/state.md T-UPSTREAM-1305-CUDA-DRAIN-BATCH-THREAD-GLOBAL-2026-09-03):
* a batch owned by another ``VmafCudaState`` must never be waited on, and
* ``vmaf_cuda_drain_batch_thread_destroy`` must leave zero entries behind so
* the next context on this thread cannot see freed events or flags. Not part
* of the public libvmaf API — ``drain_batch.h`` is a private header.
*/
unsigned vmaf_cuda_drain_batch_pending(void);

#ifdef __cplusplus
} /* extern "C" */
#endif
Expand Down
73 changes: 71 additions & 2 deletions core/src/libvmaf.c
Original file line number Diff line number Diff line change
Expand Up @@ -2666,7 +2666,7 @@ static int read_pictures_extractor_loop_cuda(VmafContext *vmaf, VmafPicture *ref
/* Phase 2: batched submit of the curr frame. The drain batch is
* re-opened so each extractor's submit() registers its
* ``finished`` event for the next frame's drain_flush. */
vmaf_cuda_drain_batch_open();
vmaf_cuda_drain_batch_open(&vmaf->cuda.state);
for (unsigned i = 0; i < vmaf->registered_feature_extractors.cnt; i++) {
VmafFeatureExtractorContext *fex_ctx = vmaf->registered_feature_extractors.fex_ctx[i];
if (!(fex_ctx->fex->flags & VMAF_FEATURE_EXTRACTOR_CUDA)) {
Expand Down Expand Up @@ -3021,6 +3021,59 @@ int vmaf_register_metadata_handler(VmafContext *vmaf, VmafMetadataConfiguration
return vmaf_feature_collector_register_metadata(vmaf->feature_collector, cfg);
}

/* Fence the pipeline so a read at ``index`` sees every write that is already
* owed for it (Netflix/vmaf#1305, docs/state.md
* T-UPSTREAM-1305-CUDA-DRAIN-BATCH-THREAD-GLOBAL-2026-09-03).
*
* The read entry points below used to go straight to the feature collector, so
* a caller pulling index N-2 while N was in flight could read a slot the
* producing thread (or the GPU collect step) had not written yet. Fencing on
* every read would serialise the pipeline, so the callers fence only after the
* lock-free read reports the slot unwritten: wait for the worker threads, then,
* on CUDA, drain the batch and run the pending collect() so the slots the
* device already computed are actually in the collector. The drain batch stays
* open afterwards — Phase 1 of the next vmaf_read_pictures() re-flushes it, and
* flushing twice is a no-op because the flush clears its entries.
*
* Returns 0 when it is worth re-reading, negative on a fence error. */
static int fence_for_read(VmafContext *vmaf, unsigned index)
{
int err = 0;

if (vmaf->thread_pool) {
err = vmaf_thread_pool_wait(vmaf->thread_pool);
if (err)
return err;
}

#ifdef HAVE_CUDA
if (vmaf->cuda.state.ctx) {
err = vmaf_cuda_drain_batch_flush(&vmaf->cuda.state);
if (err)
return err;
for (unsigned i = 0; i < vmaf->registered_feature_extractors.cnt; i++) {
VmafFeatureExtractorContext *fex_ctx = vmaf->registered_feature_extractors.fex_ctx[i];
if (!(fex_ctx->fex->flags & VMAF_FEATURE_EXTRACTOR_CUDA))
continue;
if (!fex_ctx->gpu_pending)
continue;
if (fex_ctx->gpu_pending_index > index)
continue;
err = vmaf_feature_extractor_context_collect(fex_ctx, fex_ctx->gpu_pending_index,
vmaf->feature_collector);
fex_ctx->gpu_pending = false;
fex_release_prev_ref(fex_ctx->fex);
if (err)
return err;
}
}
#else
(void)index;
#endif

return 0;
}

int vmaf_feature_score_at_index(VmafContext *vmaf, const char *feature_name, double *score,
unsigned index)
{
Expand All @@ -3031,7 +3084,16 @@ int vmaf_feature_score_at_index(VmafContext *vmaf, const char *feature_name, dou
if (!score)
return -EINVAL;

return vmaf_feature_collector_get_score(vmaf->feature_collector, feature_name, score, index);
int err = vmaf_feature_collector_get_score(vmaf->feature_collector, feature_name, score, index);
if (err == -EAGAIN) {
/* The slot exists but is unwritten: fence and read once more before
* telling the caller the frame is not ready (Netflix/vmaf#1305). */
const int fence_err = fence_for_read(vmaf, index);
if (fence_err)
return fence_err;
err = vmaf_feature_collector_get_score(vmaf->feature_collector, feature_name, score, index);
}
return err;
}

int vmaf_score_at_index(VmafContext *vmaf, VmafModel *model, double *score, unsigned index)
Expand All @@ -3057,6 +3119,13 @@ int vmaf_score_at_index(VmafContext *vmaf, VmafModel *model, double *score, unsi
* first, causing vmaf_score_pooled to return -EAGAIN for multi-frame
* sequences. ADR-1073. */
if (err) {
/* Netflix/vmaf#1305: the input features for this index may still be in
* flight (worker threads, or a CUDA collect that has not run yet), so
* fence before predicting — otherwise the prediction is computed from
* unwritten slots. */
const int fence_err = fence_for_read(vmaf, index);
if (fence_err)
return fence_err;
err = vmaf_predict_score_at_index(model, vmaf->feature_collector, index, score, true, false,
0);
}
Expand Down
15 changes: 15 additions & 0 deletions core/test/meson.build
Original file line number Diff line number Diff line change
Expand Up @@ -1612,8 +1612,23 @@ test_cuda_pic_preallocation = executable('test_cuda_pic_preallocation',
c_args: ['-DHAVE_CUDA=1']
)

# Owner scoping of the ADR-0242 drain batch
# (docs/state.md T-UPSTREAM-1305-CUDA-DRAIN-BATCH-THREAD-GLOBAL-2026-09-03).
# Exercises only the bookkeeping paths — a foreign owner short-circuits the
# flush and the destroy path returns as soon as it sees a NULL drain stream —
# so it needs nvcc to compile but no live device, and is deliberately NOT in
# the 'gpu' suite.
test_cuda_drain_batch = executable('test_cuda_drain_batch',
['test.c', 'test_cuda_drain_batch.c'],
include_directories : [libvmaf_inc, test_inc, include_directories('../src/')],
link_with : get_option('default_library') == 'both' ? libvmaf.get_static_lib() : libvmaf,
dependencies: [pthread_dependency, cuda_dependency],
c_args: ['-DHAVE_CUDA=1']
)

test('test_gpu_picture_pool', test_gpu_picture_pool, suite : ['fast', 'gpu'])
test('test_cuda_pic_preallocation', test_cuda_pic_preallocation, suite : ['fast', 'gpu'])
test('test_cuda_drain_batch', test_cuda_drain_batch, suite : ['fast'])

# Finding R2-2 — HOST_PINNED picture-allocator dimension-overflow guard.
# vmaf_cuda_picture_alloc_pinned now rejects w/h out of [1, 32768] BEFORE
Expand Down
Loading
Loading