Skip to content
Closed
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 @@ -9014,6 +9014,12 @@ non-empty, printable string. Closes a coverage gap noted in
`fuzz/meson.build`; registered in suite `fast`.


- Allow `tidy-ratchet.py --only <TU> --write` to tighten measured source
allowances while preserving unmeasured entries and full-report metadata.
Exact coverage, matching tool version, parse/compile success, debt monotonicity
and atomic replacement are required; CI still measures the whole tree.


- **`--tiny-codec` / `--tiny-preset` / `--tiny-crf` CLI flags** populate
the codec one-hot block of codec-aware tiny models (today
`fr_regressor_v2`) so the model receives the real encoder context
Expand Down Expand Up @@ -26285,6 +26291,13 @@ ADR-0513.
- Test YUV fixture provisioner: `scripts/test/fetch-test-yuvs.sh` downloads `src01_hrc0[0-1]_576x324.yuv` from `Netflix/vmaf_resource` and md5-verifies them. Reverts the [#1237](https://github.com/VMAFx/vmafx/pull/1237) ADM2 golden override, which was based on output from stale local fixture content. See [ADR-0493](docs/adr/0493-test-yuv-fixture-md5-verification.md) and [docs/development/test-fixtures.md](docs/development/test-fixtures.md).


- Bound pending CPU thread-pool jobs to the created worker count, restoring
backpressure when decoding outpaces feature extraction. Preserve recycled
payloads and batch errors, and wait for blocked producers during shutdown.
Adapted from Netflix/vmaf commit `8fc71e3` with shutdown-lifetime regression
coverage.


- `vmaf_thread_pool_create` now checks the return value of `pthread_create`
and handles partial-success (at least one thread started) and total-failure
(zero threads started, returns `-EAGAIN`/`-EPERM` to the caller) gracefully.
Expand Down
4 changes: 4 additions & 0 deletions changelog.d/added/tidy-scoped-baseline-tightening.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
- Allow `tidy-ratchet.py --only <TU> --write` to tighten measured source
allowances while preserving unmeasured entries and full-report metadata.
Exact coverage, matching tool version, parse/compile success, debt monotonicity
and atomic replacement are required; CI still measures the whole tree.
5 changes: 5 additions & 0 deletions changelog.d/fixed/thread-pool-bounded-queue.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
- Bound pending CPU thread-pool jobs to the created worker count, restoring
backpressure when decoding outpaces feature extraction. Preserve recycled
payloads and batch errors, and wait for blocked producers during shutdown.
Adapted from Netflix/vmaf commit `8fc71e3` with shutdown-lifetime regression
coverage.
10 changes: 10 additions & 0 deletions core/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -156,6 +156,16 @@ core/
`boolean` (matching `enable_cuda` / `enable_sycl`); do NOT
convert it to `feature` without an ADR amendment per ADR-0212
§ "Decision".
- **Bounded thread-pool admission** (Netflix `8fc71e3`, fork lifetime adaptation):
`src/thread_pool.c` admits at most one queued job per successfully created
worker, before payload allocation. Keep dequeue wakeups and checked condition
teardown. Destruction waits for both workers and already-blocked producers;
waking a producer does not make it safe to free its mutex. Preserve the
cancellation/lifetime and mixed-payload tests in
`test/test_thread_pool_backpressure.c`. Callbacks must not enqueue into the
same pool. Serialize destruction against API entry, including mutex-acquisition
waiters; only proven registered capacity waiters can be cancelled concurrently. See
[thread-pool behavior](../docs/development/thread-pool.md).
- **Thread-pool job recycling + inline data buffer** (fork-local,
ADR-0147): [`src/thread_pool.c`](src/thread_pool.c) recycles
`VmafThreadPoolJob` slots via a `pool->free_jobs` free list
Expand Down
67 changes: 60 additions & 7 deletions core/src/thread_pool.c
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,9 @@

#include "thread_pool.h"

/* MSVC C23 requires NULL in C sources (ADR-1138). */
// NOLINTBEGIN(modernize-use-nullptr)

/* Payload ≤ JOB_INLINE_DATA_SIZE lives inside the job struct itself,
* avoiding a second malloc per enqueue. Sized to cover every
* `data_sz` used by current callers (the CPU read_pictures stage
Expand All @@ -47,7 +50,10 @@ typedef struct VmafThreadPool {
struct {
pthread_mutex_t lock;
pthread_cond_t empty;
pthread_cond_t not_full;
VmafThreadPoolJob *head, *tail;
unsigned depth;
unsigned waiting_producers;
} queue;
pthread_cond_t working;
/* n_threads: live worker count; decremented by each runner on exit (lock held).
Expand Down Expand Up @@ -129,6 +135,10 @@ static void *vmaf_thread_pool_runner(void *p)
if (pool->stop)
break;
VmafThreadPoolJob *job = vmaf_thread_pool_fetch_job(pool);
if (job) {
pool->queue.depth--;
(void)pthread_cond_signal(&pool->queue.not_full);
}
pool->n_working++;
pthread_mutex_unlock(&(pool->queue.lock));
int job_err = 0;
Expand All @@ -153,8 +163,8 @@ static void *vmaf_thread_pool_runner(void *p)
return NULL;
}

/* Initialise the three synchronisation primitives in `p` in dependency
* order (mutex first, then both cond vars). On failure, tears down only
/* Initialise the synchronisation primitives in `p` in dependency
* order (mutex first, then the cond vars). On failure, tears down only
* the objects that were successfully initialised and returns -ENOMEM.
* pthread_*_init can fail with ENOMEM on constrained / embedded systems;
* ignoring the return value leaves the pool in undefined state. */
Expand All @@ -173,7 +183,16 @@ static int pool_init_primitives(VmafThreadPool *p, VmafThreadPool **pool_out_to_
*pool_out_to_null = NULL;
return -ENOMEM;
}
if (pthread_cond_init(&(p->queue.not_full), NULL) != 0) {
pthread_cond_destroy(&(p->queue.empty));
pthread_mutex_destroy(&(p->queue.lock));
free(p->workers);
free(p);
*pool_out_to_null = NULL;
return -ENOMEM;
}
if (pthread_cond_init(&(p->working), NULL) != 0) {
pthread_cond_destroy(&(p->queue.not_full));
pthread_cond_destroy(&(p->queue.empty));
pthread_mutex_destroy(&(p->queue.lock));
free(p->workers);
Expand All @@ -200,6 +219,7 @@ static int pool_spawn_workers(VmafThreadPool *p, VmafThreadPoolConfig cfg,
if (i == 0) {
pthread_mutex_destroy(&(p->queue.lock));
pthread_cond_destroy(&(p->queue.empty));
pthread_cond_destroy(&(p->queue.not_full));
pthread_cond_destroy(&(p->working));
free(p->workers);
free(p);
Expand Down Expand Up @@ -244,8 +264,28 @@ int vmaf_thread_pool_create(VmafThreadPool **pool, VmafThreadPoolConfig cfg)
return pool_spawn_workers(p, cfg, pool);
}

/* Caller holds queue.lock throughout admission and job allocation. */
static int wait_for_queue_capacity(VmafThreadPool *pool)
{
/* Bound queued work before allocating or retaining a job payload. The
* actual created width also covers partial pthread_create failure. */
pool->queue.waiting_producers++;
while (pool->queue.depth >= pool->n_workers_created && !pool->stop)
(void)pthread_cond_wait(&pool->queue.not_full, &pool->queue.lock);
pool->queue.waiting_producers--;
if (pool->stop) {
/* Destroy must not free the condition/mutex while a producer is still
* waking from its capacity wait. API entry is externally serialized
* against destroy; callers waiting to acquire queue.lock are not counted. */
(void)pthread_cond_signal(&pool->working);
return -ECANCELED;
}

return 0;
}

int vmaf_thread_pool_enqueue(VmafThreadPool *pool, int (*func)(void *data, void **thread_data),
void *data, size_t data_sz)
const void *data, size_t data_sz)
{
if (!pool)
return -EINVAL;
Expand All @@ -254,9 +294,14 @@ int vmaf_thread_pool_enqueue(VmafThreadPool *pool, int (*func)(void *data, void

pthread_mutex_lock(&(pool->queue.lock));

/* Reuse a recycled slot if available, otherwise heap-allocate one.
* The free list is mutex-protected so fetching and returning stays
* coherent with the runner thread's recycle on job completion. */
const int admission_err = wait_for_queue_capacity(pool);
if (admission_err) {
(void)pthread_mutex_unlock(&pool->queue.lock);
return admission_err;
}

/* The mutex protects slot reuse against worker-side recycling; allocate
* a slot only when none is free. */
VmafThreadPoolJob *job = pool->free_jobs;
if (job) {
pool->free_jobs = job->next;
Expand Down Expand Up @@ -296,6 +341,7 @@ int vmaf_thread_pool_enqueue(VmafThreadPool *pool, int (*func)(void *data, void
pool->queue.tail = job;
}

pool->queue.depth++;
pthread_cond_signal(&(pool->queue.empty));
pthread_mutex_unlock(&(pool->queue.lock));

Expand All @@ -309,7 +355,7 @@ int vmaf_thread_pool_wait(VmafThreadPool *pool)

pthread_mutex_lock(&(pool->queue.lock));
while ((!pool->stop && (pool->n_working || pool->queue.head)) ||
(pool->stop && pool->n_threads))
(pool->stop && (pool->n_threads || pool->queue.waiting_producers)))
pthread_cond_wait(&(pool->working), &(pool->queue.lock));
/* Harvest and clear the accumulated worker-error flags. Clearing here
* means each vmaf_thread_pool_wait() call reports errors from the batch
Expand Down Expand Up @@ -339,8 +385,12 @@ int vmaf_thread_pool_destroy(VmafThreadPool *pool)
job = next_job;
}

pool->queue.head = NULL;
pool->queue.tail = NULL;
pool->queue.depth = 0;
pool->stop = true;
pthread_cond_broadcast(&(pool->queue.empty));
(void)pthread_cond_broadcast(&pool->queue.not_full);
pthread_mutex_unlock(&(pool->queue.lock));
vmaf_thread_pool_wait(pool);

Expand All @@ -364,8 +414,11 @@ int vmaf_thread_pool_destroy(VmafThreadPool *pool)

pthread_mutex_destroy(&(pool->queue.lock));
pthread_cond_destroy(&(pool->queue.empty));
(void)pthread_cond_destroy(&(pool->queue.not_full));
pthread_cond_destroy(&(pool->working));

free(pool);
return 0;
}

// NOLINTEND(modernize-use-nullptr)
21 changes: 16 additions & 5 deletions core/src/thread_pool.h
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,13 @@ typedef struct VmafThreadPoolConfig {
int vmaf_thread_pool_create(VmafThreadPool **tpool, VmafThreadPoolConfig cfg);

/**
* @brief Enqueue a work item on the pool.
* @brief Enqueue a work item, waiting while the bounded queue is full.
*
* At most one pending job per successfully created worker is retained. A
* dequeuing worker wakes a producer; running jobs are outside this bound.
* A registered capacity waiter returns -ECANCELED during coordinated shutdown.
* Callbacks must not recursively enqueue into this same pool or wait on work
* whose submission requires the blocked producer to advance.
*
* @p data is copied internally (up to @p data_sz bytes), so the caller's
* buffer may be reused or freed immediately after this call returns.
Expand All @@ -68,7 +74,7 @@ int vmaf_thread_pool_create(VmafThreadPool **tpool, VmafThreadPoolConfig cfg);
* @return 0 on success, negative errno on failure.
*/
int vmaf_thread_pool_enqueue(VmafThreadPool *pool, int (*func)(void *data, void **thread_data),
void *data, size_t data_sz);
const void *data, size_t data_sz);

/**
* @brief Block until all enqueued work items have completed.
Expand All @@ -79,11 +85,16 @@ int vmaf_thread_pool_enqueue(VmafThreadPool *pool, int (*func)(void *data, void
int vmaf_thread_pool_wait(VmafThreadPool *pool);

/**
* @brief Drain the pool and free all associated resources.
* @brief Discard queued jobs, finish active jobs, and free the pool.
*
* Implicitly calls vmaf_thread_pool_wait() before tearing down threads.
* Call vmaf_thread_pool_wait() first to drain all accepted work. Destruction
* wakes and waits for registered queue-capacity waiters. The owner must
* externally serialize destruction against API entry, including callers
* waiting to acquire the queue mutex. Cancellation is supported only when
* the owner has established that the concurrent producers are already in
* the capacity wait. Otherwise stop and join producers before destruction.
*
* @param tpool Pool to destroy. May be NULL (no-op).
* @param tpool Pool to destroy. NULL returns -EINVAL.
* @return 0 on success, negative errno on failure.
*/
int vmaf_thread_pool_destroy(VmafThreadPool *tpool);
Expand Down
8 changes: 8 additions & 0 deletions core/test/meson.build
Original file line number Diff line number Diff line change
Expand Up @@ -198,6 +198,13 @@ test_thread_pool = executable('test_thread_pool',
dependencies : [pthread_dependency, thread_lib, sycl_dependency],
)

# Source inclusion allows deterministic pthread failure and wakeup injection.
test_thread_pool_backpressure = executable('test_thread_pool_backpressure',
['test.c', 'test_thread_pool_backpressure.c'],
include_directories : [libvmaf_inc, test_inc, include_directories('../src/')],
dependencies : [pthread_dependency, thread_lib],
)

test_model = executable('test_model',
['test.c', 'test_model.c', '../src/dict.cpp', '../src/pdjson.c', json_model_c_sources],
include_directories : [libvmaf_inc, test_inc, include_directories('../src')],
Expand Down Expand Up @@ -3525,6 +3532,7 @@ test('test_output', test_output, suite : ['fast'])
test('test_log', test_log, suite : ['fast'])
test('test_public_api_score', test_public_api_score, suite : ['fast'])
test('test_thread_pool', test_thread_pool, suite : ['fast'])
test('test_thread_pool_backpressure', test_thread_pool_backpressure, suite : ['fast'])
test('test_model', test_model, suite : ['fast'])
test('test_model_libsvm_dup_key', test_model_libsvm_dup_key, suite : ['fast'])
test('test_model_feature_overload_ownership', test_model_feature_overload_ownership, suite : ['fast'])
Expand Down
Loading
Loading