Browse Source

Deliver media completions without polling or timer failure joins

master
evgeny 2 days ago
parent
commit
adff64e483
  1. 34
      src/media_async/media_async.c
  2. 5
      src/media_async/media_async.h

34
src/media_async/media_async.c

@ -27,12 +27,11 @@ struct ma_thread_ctx {
struct ma_thread_ctx* next;
struct UASYNC* ua;
pthread_t thread;
void* timer;
struct posted_task* completion;
ma_work_fn work;
void* data;
ma_done_fn done;
void* arg;
int finished;
};
static void ma_finish(struct ma_thread_ctx* job, int error) {
@ -44,21 +43,17 @@ static void ma_finish(struct ma_thread_ctx* job, int error) {
u_free(job);
}
static void ma_poll(void* arg) {
struct ma_thread_ctx* job = arg; job->timer = NULL;
if (__atomic_load_n(&job->finished, __ATOMIC_ACQUIRE)) { ma_finish(job, 0); return; }
job->timer = uasync_set_timeout(job->ua, 100, job, ma_poll, "media_worker");
if (!job->timer) {
DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "media worker completion timer failed");
ma_finish(job, -1);
}
/* Completion приходит только после работы; ожидать незавершённый worker здесь нельзя. */
static void ma_completed(void* arg) {
struct ma_thread_ctx* job = arg;
job->completion = NULL;
ma_finish(job, 0);
}
static void* ma_thread_entry(void* arg) {
struct ma_thread_ctx* job = arg;
job->work(job->data);
__atomic_store_n(&job->finished, 1, __ATOMIC_RELEASE);
uasync_wakeup(job->ua);
uasync_post_reserved(job->ua, job->completion);
return NULL;
}
@ -75,8 +70,11 @@ void media_async_destroy(struct media_async* ma) {
unsigned cancelled = 0;
while (ma->jobs) {
struct ma_thread_ctx* job = ma->jobs;
if (job->timer) uasync_cancel_timeout(job->ua, job->timer);
ma_finish(job, MEDIA_ASYNC_CANCELLED); cancelled++;
pthread_join(job->thread, NULL);
uasync_cancel_post(job->ua, job->completion);
ma->jobs = job->next;
job->done(job->arg, MEDIA_ASYNC_CANCELLED);
u_free(job); cancelled++;
}
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "media_async: destroyed cancelled=%u", cancelled);
u_free(ma);
@ -92,12 +90,14 @@ void media_async_submit(struct media_async* ma, struct UASYNC* ua,
struct ma_thread_ctx* job = u_calloc(1, sizeof(*job));
if (!job) { DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "media task allocation failed"); done(arg, -1); return; }
job->owner = ma; job->ua = ua; job->work = work; job->data = data; job->done = done; job->arg = arg;
job->timer = uasync_set_timeout(ua, 1, job, ma_poll, "media_worker");
if (!job->timer) { DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "media task timer failed"); u_free(job); done(arg, -1); return; }
job->completion = u_calloc(1, sizeof(*job->completion));
if (!job->completion) { DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "media completion allocation failed"); u_free(job); done(arg, -1); return; }
job->completion->callback = ma_completed;
job->completion->arg = job;
int rc = pthread_create(&job->thread, NULL, ma_thread_entry, job);
if (rc) {
DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "media task thread failed rc=%d", rc);
uasync_cancel_timeout(ua, job->timer); u_free(job); done(arg, -1); return;
u_free(job->completion); u_free(job); done(arg, -1); return;
}
job->next = ma->jobs; ma->jobs = job;
}

5
src/media_async/media_async.h

@ -17,12 +17,13 @@ struct UASYNC;
struct media_async;
typedef void (*ma_work_fn)(void* data); // рабочий поток; не обращаться к состоянию uasync
typedef void (*ma_done_fn)(void* arg, int err); // 0=work завершён, -1=ошибка запуска/таймера, -2=отмена
typedef void (*ma_done_fn)(void* arg, int err); // 0=work завершён, -1=ошибка запуска, -2=отмена
/* Создать владельца задач; NULL при ошибке выделения памяти. */
struct media_async* media_async_create(void);
#define MEDIA_ASYNC_CANCELLED (-2)
/* Вызывать вне done. Ждёт работы, вызывает ожидающие done и освобождает ma; NULL допустим. */
/* Вызывать вне callbacks uasync. Ждёт работы, отменяет posted completion,
* вызывает ожидающие done и освобождает ma; NULL допустим. */
void media_async_destroy(struct media_async* ma);
/* Передать задачу; data/arg остаются у вызывающего до done. При ошибке done

Loading…
Cancel
Save