Asterisk - The Open Source Telephony Project GIT-master-70eff7f
Loading...
Searching...
No Matches
Data Structures | Macros | Typedefs | Enumerations | Functions | Variables
taskpool.c File Reference
#include "asterisk.h"
#include "asterisk/_private.h"
#include "asterisk/taskpool.h"
#include "asterisk/taskprocessor.h"
#include "asterisk/astobj2.h"
#include "asterisk/serializer_shutdown_group.h"
#include "asterisk/utils.h"
#include "asterisk/time.h"
#include "asterisk/sched.h"
Include dependency graph for taskpool.c:

Go to the source code of this file.

Data Structures

struct  ast_taskpool
 An opaque taskpool structure. More...
 
struct  serializer
 
struct  taskpool_sync_task
 
struct  taskpool_taskprocessor
 A taskpool taskprocessor. More...
 
struct  taskpool_taskprocessors
 A container of taskprocessors. More...
 

Macros

#define ast_taskpool_push_internal(pool, task, data)    __ast_taskpool_push(pool, task, data, __FILE__, __LINE__, __PRETTY_FUNCTION__)
 
#define TASKPOOL_GROW_THRESHOLD   (AST_TASKPROCESSOR_HIGH_WATER_LEVEL * 5) / 10
 The threshold for a taskprocessor at which we consider the pool needing to grow (50% of high water threshold)
 
#define TASKPOOL_QUEUE_SIZE_ADD(tps, size)   (size += ast_taskprocessor_size(tps->taskprocessor))
 
#define TASKPROCESSOR_IS_IDLE(tps, timeout)   (ast_tvdiff_ms(ast_tvnow(), tps->last_pushed) > (timeout))
 

Typedefs

typedef void(* taskpool_selector) (struct ast_taskpool *pool, struct taskpool_taskprocessors *taskprocessors, struct taskpool_taskprocessor **taskprocessor, unsigned int *growth_threshold_reached)
 

Enumerations

enum  serializer_suspension { SERIALIZER_UNSUSPENDED = 0 , SERIALIZER_SUSPENDING , SERIALIZER_SUSPENDED }
 

Functions

int __ast_taskpool_push (struct ast_taskpool *pool, int(*task)(void *data), void *data, const char *file, int line, const char *function)
 Push a task to the taskpool.
 
int __ast_taskpool_push_wait (struct ast_taskpool *pool, int(*task)(void *data), void *data, const char *file, int line, const char *function)
 Push a task to the taskpool, and wait for completion.
 
int __ast_taskpool_serializer_push_wait (struct ast_taskprocessor *serializer, int(*task)(void *data), void *data, const char *file, int line, const char *function)
 Push a task to a serializer, and wait for completion.
 
struct ast_taskpoolast_taskpool_create (const char *name, const struct ast_taskpool_options *options)
 Create a new taskpool.
 
static struct ast_taskpoolast_taskpool_get_current (void)
 
int ast_taskpool_init (void)
 
int ast_taskpool_push (struct ast_taskpool *pool, int(*task)(void *data), void *data)
 
int ast_taskpool_push_wait (struct ast_taskpool *pool, int(*task)(void *data), void *data)
 
long ast_taskpool_queue_size (struct ast_taskpool *pool)
 Get the current number of queued tasks in the taskpool.
 
struct ast_taskprocessorast_taskpool_serializer (const char *name, struct ast_taskpool *pool)
 Serialized execution of tasks within a ast_taskpool.
 
struct ast_taskprocessorast_taskpool_serializer_get_current (void)
 Get the taskpool serializer currently associated with this thread.
 
struct ast_taskprocessorast_taskpool_serializer_group (const char *name, struct ast_taskpool *pool, struct ast_serializer_shutdown_group *shutdown_group)
 Serialized execution of tasks within a ast_taskpool.
 
int ast_taskpool_serializer_push_wait (struct ast_taskprocessor *serializer, int(*task)(void *data), void *data)
 
int ast_taskpool_serializer_suspend (struct ast_taskprocessor *serializer)
 Suspend a serializer, causing tasks to be queued until unsuspended.
 
int ast_taskpool_serializer_unsuspend (struct ast_taskprocessor *serializer)
 Unsuspend a serializer, causing tasks to be executed.
 
void ast_taskpool_shutdown (struct ast_taskpool *pool)
 Shut down a taskpool and remove the underlying taskprocessors.
 
size_t ast_taskpool_taskprocessors_count (struct ast_taskpool *pool)
 Get the current number of taskprocessors in the taskpool.
 
 AST_THREADSTORAGE_RAW (current_taskpool_pool)
 Thread storage for the current taskpool.
 
 AST_THREADSTORAGE_RAW (current_taskpool_serializer)
 
static int execute_tasks (void *data)
 
static struct serializerserializer_create (struct ast_taskpool *pool, struct ast_serializer_shutdown_group *shutdown_group)
 
static void serializer_dtor (void *obj)
 
static void serializer_shutdown (struct ast_taskprocessor_listener *listener)
 
static int serializer_start (struct ast_taskprocessor_listener *listener)
 
static void serializer_task_pushed (struct ast_taskprocessor_listener *listener, int was_empty)
 
static void taskpool_dynamic_pool_grow (struct ast_taskpool *pool, struct taskpool_taskprocessor **taskprocessor)
 
static int taskpool_dynamic_pool_shrink (const void *data)
 
static void taskpool_least_full_selector (struct ast_taskpool *pool, struct taskpool_taskprocessors *taskprocessors, struct taskpool_taskprocessor **taskprocessor, unsigned int *growth_threshold_reached)
 Least full taskprocessor selector.
 
static void taskpool_sequential_selector (struct ast_taskpool *pool, struct taskpool_taskprocessors *taskprocessors, struct taskpool_taskprocessor **taskprocessor, unsigned int *growth_threshold_reached)
 
static int taskpool_serializer_empty_task (void *data)
 
static int taskpool_serializer_suspend_task (void *data)
 
static void taskpool_shutdown (void)
 
static int taskpool_sync_task (void *data)
 
static void taskpool_sync_task_cleanup (struct taskpool_sync_task *sync_task)
 
static int taskpool_sync_task_init (struct taskpool_sync_task *sync_task, int(*task)(void *), void *data)
 
static struct taskpool_taskprocessortaskpool_taskprocessor_alloc (struct ast_taskpool *pool, char type)
 
static void taskpool_taskprocessor_dtor (void *obj)
 
static long taskpool_taskprocessor_load (struct taskpool_taskprocessor *tp)
 
static int taskpool_taskprocessor_start (void *data)
 
static int taskpool_taskprocessor_stop (void *data)
 
static void taskpool_taskprocessors_cleanup (struct taskpool_taskprocessors *taskprocessors)
 
static int taskpool_taskprocessors_init (struct taskpool_taskprocessors *taskprocessors, unsigned int size)
 

Variables

static struct ast_sched_contextsched
 Scheduler used for dynamic pool shrinking.
 
static struct ast_taskprocessor_listener_callbacks serializer_tps_listener_callbacks
 

Macro Definition Documentation

◆ ast_taskpool_push_internal

#define ast_taskpool_push_internal (   pool,
  task,
  data 
)     __ast_taskpool_push(pool, task, data, __FILE__, __LINE__, __PRETTY_FUNCTION__)

Definition at line 540 of file taskpool.c.

◆ TASKPOOL_GROW_THRESHOLD

#define TASKPOOL_GROW_THRESHOLD   (AST_TASKPROCESSOR_HIGH_WATER_LEVEL * 5) / 10

The threshold for a taskprocessor at which we consider the pool needing to grow (50% of high water threshold)

Definition at line 80 of file taskpool.c.

◆ TASKPOOL_QUEUE_SIZE_ADD

#define TASKPOOL_QUEUE_SIZE_ADD (   tps,
  size 
)    (size += ast_taskprocessor_size(tps->taskprocessor))

Definition at line 479 of file taskpool.c.

◆ TASKPROCESSOR_IS_IDLE

#define TASKPROCESSOR_IS_IDLE (   tps,
  timeout 
)    (ast_tvdiff_ms(ast_tvnow(), tps->last_pushed) > (timeout))

Definition at line 235 of file taskpool.c.

Typedef Documentation

◆ taskpool_selector

typedef void(* taskpool_selector) (struct ast_taskpool *pool, struct taskpool_taskprocessors *taskprocessors, struct taskpool_taskprocessor **taskprocessor, unsigned int *growth_threshold_reached)

Definition at line 51 of file taskpool.c.

Enumeration Type Documentation

◆ serializer_suspension

Enumerator
SERIALIZER_UNSUSPENDED 
SERIALIZER_SUSPENDING 
SERIALIZER_SUSPENDED 

Definition at line 713 of file taskpool.c.

713 {
714 SERIALIZER_UNSUSPENDED = 0, /* The serializer is unsuspended */
715 SERIALIZER_SUSPENDING, /* The serializer is pending suspension */
716 SERIALIZER_SUSPENDED, /* The serializer is suspended */
717};
@ SERIALIZER_SUSPENDED
Definition taskpool.c:716
@ SERIALIZER_SUSPENDING
Definition taskpool.c:715
@ SERIALIZER_UNSUSPENDED
Definition taskpool.c:714

Function Documentation

◆ __ast_taskpool_push()

int __ast_taskpool_push ( struct ast_taskpool pool,
int(*)(void *data)  task,
void *  data,
const char *  file,
int  line,
const char *  function 
)

Push a task to the taskpool.

Since
23.1.0
22.7.0
20.17.0

Tasks pushed into the taskpool will be automatically taken by one of the taskprocessors within

Parameters
poolThe taskpool to add the task to
taskThe task to add
dataThe parameter for the task
Return values
0success
-1failure

Definition at line 544 of file taskpool.c.

546{
547 RAII_VAR(struct taskpool_taskprocessor *, taskprocessor, NULL, ao2_cleanup);
548
549 /* Select the taskprocessor in the pool to use for pushing this task */
550 ao2_lock(pool);
551 if (!pool->shutting_down) {
552 unsigned int growth_threshold_reached = 0;
553
554 /* A selector doesn't set taskprocessor to NULL, it will only change the value if a better
555 * taskprocessor is found. This means that even if the selector for a dynamic taskprocessor
556 * fails for some reason, it will still fall back to the initially found static one if
557 * it is present.
558 */
559 pool->selector(pool, &pool->static_taskprocessors, &taskprocessor, &growth_threshold_reached);
560 if (pool->options.auto_increment && growth_threshold_reached) {
561 /* If we need to grow then try dynamic taskprocessors */
562 pool->selector(pool, &pool->dynamic_taskprocessors, &taskprocessor, &growth_threshold_reached);
563 if (growth_threshold_reached) {
564 /* If we STILL need to grow then grow the dynamic taskprocessor pool if allowed */
565 taskpool_dynamic_pool_grow(pool, &taskprocessor);
566 }
567
568 /* If a dynamic taskprocessor was used update its last push time */
569 if (taskprocessor) {
570 taskprocessor->last_pushed = ast_tvnow();
571 }
572 }
573 ao2_bump(taskprocessor);
574 }
575 ao2_unlock(pool);
576
577 if (!taskprocessor) {
578 return -1;
579 }
580
581 if (__ast_taskprocessor_push(taskprocessor->taskprocessor, task, data, file, line, function)) {
582 return -1;
583 }
584
585 return 0;
586}
#define ao2_cleanup(obj)
Definition astobj2.h:1934
#define ao2_unlock(a)
Definition astobj2.h:729
#define ao2_lock(a)
Definition astobj2.h:717
#define ao2_bump(obj)
Bump refcount on an AO2 object by one, returning the object.
Definition astobj2.h:480
#define NULL
Definition resample.c:96
int auto_increment
Number of taskprocessors to increment the pool by.
Definition taskpool.h:92
taskpool_selector selector
Definition taskpool.c:74
struct ast_taskpool_options options
Definition taskpool.c:70
struct taskpool_taskprocessors static_taskprocessors
Definition taskpool.c:64
int shutting_down
Definition taskpool.c:68
struct taskpool_taskprocessors dynamic_taskprocessors
Definition taskpool.c:66
A taskpool taskprocessor.
Definition taskpool.c:34
static void taskpool_dynamic_pool_grow(struct ast_taskpool *pool, struct taskpool_taskprocessor **taskprocessor)
Definition taskpool.c:496
int __ast_taskprocessor_push(struct ast_taskprocessor *tps, int(*task_exe)(void *datap), void *datap, const char *file, int line, const char *function) attribute_warn_unused_result
Push a task into the specified taskprocessor queue and signal the taskprocessor thread.
static int task(void *data)
Queued task for baseline test.
struct timeval ast_tvnow(void)
Returns current timeval. Meant to replace calls to gettimeofday().
Definition time.h:159
#define RAII_VAR(vartype, varname, initval, dtor)
Declare a variable that will call a destructor function when it goes out of scope.
Definition utils.h:981

References __ast_taskprocessor_push(), ao2_bump, ao2_cleanup, ao2_lock, ao2_unlock, ast_tvnow(), ast_taskpool_options::auto_increment, ast_taskpool::dynamic_taskprocessors, NULL, ast_taskpool::options, RAII_VAR, ast_taskpool::selector, ast_taskpool::shutting_down, ast_taskpool::static_taskprocessors, task(), taskpool_dynamic_pool_grow(), and taskpool_taskprocessor::taskprocessor.

Referenced by __ast_sip_push_task(), __ast_taskpool_push_wait(), and ast_taskpool_push().

◆ __ast_taskpool_push_wait()

int __ast_taskpool_push_wait ( struct ast_taskpool pool,
int(*)(void *data)  task,
void *  data,
const char *  file,
int  line,
const char *  function 
)

Push a task to the taskpool, and wait for completion.

Since
23.1.0
22.7.0
20.17.0

Tasks pushed into the taskpool will be automatically taken by one of the taskprocessors within

Parameters
poolThe taskpool to add the task to
taskThe task to add
dataThe parameter for the task
Return values
0success
-1failure

Definition at line 652 of file taskpool.c.

654{
655 struct taskpool_sync_task sync_task;
656
657 /* If we are already executing within a taskpool taskprocessor then
658 * don't bother pushing a new task, just directly execute the task.
659 */
661 return task(data);
662 }
663
664 if (taskpool_sync_task_init(&sync_task, task, data)) {
665 return -1;
666 }
667
668 if (__ast_taskpool_push(pool, taskpool_sync_task, &sync_task, file, line, function)) {
669 taskpool_sync_task_cleanup(&sync_task);
670 return -1;
671 }
672
673 ast_mutex_lock(&sync_task.lock);
674 while (!sync_task.complete) {
675 ast_cond_wait(&sync_task.cond, &sync_task.lock);
676 }
677 ast_mutex_unlock(&sync_task.lock);
678
679 taskpool_sync_task_cleanup(&sync_task);
680 return sync_task.fail;
681}
#define ast_cond_wait(cond, mutex)
Definition lock.h:212
#define ast_mutex_unlock(a)
Definition lock.h:197
#define ast_mutex_lock(a)
Definition lock.h:196
static void taskpool_sync_task_cleanup(struct taskpool_sync_task *sync_task)
Definition taskpool.c:623
int __ast_taskpool_push(struct ast_taskpool *pool, int(*task)(void *data), void *data, const char *file, int line, const char *function)
Push a task to the taskpool.
Definition taskpool.c:544
static struct ast_taskpool * ast_taskpool_get_current(void)
Definition taskpool.c:106
static int taskpool_sync_task_init(struct taskpool_sync_task *sync_task, int(*task)(void *), void *data)
Definition taskpool.c:609

References __ast_taskpool_push(), ast_cond_wait, ast_mutex_lock, ast_mutex_unlock, ast_taskpool_get_current(), taskpool_sync_task::complete, taskpool_sync_task::cond, taskpool_sync_task::fail, taskpool_sync_task::lock, task(), taskpool_sync_task_cleanup(), and taskpool_sync_task_init().

Referenced by __ast_sip_push_task_wait(), __ast_sip_push_task_wait_serializer(), and ast_taskpool_push_wait().

◆ __ast_taskpool_serializer_push_wait()

int __ast_taskpool_serializer_push_wait ( struct ast_taskprocessor serializer,
int(*)(void *data)  task,
void *  data,
const char *  file,
int  line,
const char *  function 
)

Push a task to a serializer, and wait for completion.

Since
23.1.0
22.7.0
20.17.0
Parameters
serializerThe serializer to add the task to
taskThe task to add
dataThe parameter for the task
Return values
0success
-1failure

Definition at line 899 of file taskpool.c.

901{
904 struct ast_taskprocessor *prior_serializer;
905 struct taskpool_sync_task sync_task;
906
907 /* If not in a taskpool taskprocessor we can just queue the task like normal and
908 * wait. */
910 if (taskpool_sync_task_init(&sync_task, task, data)) {
911 return -1;
912 }
913
914 if (__ast_taskprocessor_push(serializer, taskpool_sync_task, &sync_task, file, line, function)) {
915 taskpool_sync_task_cleanup(&sync_task);
916 return -1;
917 }
918
919 ast_mutex_lock(&sync_task.lock);
920 while (!sync_task.complete) {
921 ast_cond_wait(&sync_task.cond, &sync_task.lock);
922 }
923 ast_mutex_unlock(&sync_task.lock);
924
925 taskpool_sync_task_cleanup(&sync_task);
926 return sync_task.fail;
927 }
928
929 /* It is possible that we are already executing within a serializer, so stash the existing
930 * away so we can restore it.
931 */
932 prior_serializer = ast_taskpool_serializer_get_current();
933
934 ao2_lock(ser);
935
936 /* There are two cases where we can or have to directly execute this task:
937 * 1. There are no other tasks in the serializer
938 * 2. We are already in the serializer
939 * In the second case if we don't execute the task now, we will deadlock waiting
940 * on it as it will never occur.
941 */
942 if (!ast_taskprocessor_size(serializer) || prior_serializer == serializer) {
943 ast_threadstorage_set_ptr(&current_taskpool_serializer, serializer);
944 sync_task.fail = task(data);
945 ao2_unlock(ser);
946 ast_threadstorage_set_ptr(&current_taskpool_serializer, prior_serializer);
947 return sync_task.fail;
948 }
949
950 if (taskpool_sync_task_init(&sync_task, task, data)) {
951 ao2_unlock(ser);
952 return -1;
953 }
954
955 /* First we queue the serialized task */
956 if (__ast_taskprocessor_push(serializer, taskpool_sync_task, &sync_task, file, line, function)) {
957 taskpool_sync_task_cleanup(&sync_task);
958 ao2_unlock(ser);
959 return -1;
960 }
961
962 /* Next we queue the empty task to ensure the serializer doesn't reach empty, this
963 * stops two tasks from being queued for the same serializer at the same time.
964 */
966 taskpool_sync_task_cleanup(&sync_task);
967 ao2_unlock(ser);
968 return -1;
969 }
970
971 /* Now we execute the tasks on the serializer until our sync task is complete */
972 ast_threadstorage_set_ptr(&current_taskpool_serializer, serializer);
973 while (!sync_task.complete) {
974 /* If the serializer is suspended wait until it unsuspends */
975 while (ser->suspended == SERIALIZER_SUSPENDED) {
977 }
978
979 /* The sync task is guaranteed to be executed, so doing a while loop on the complete
980 * flag is safe.
981 */
983 }
984 taskpool_sync_task_cleanup(&sync_task);
985 ao2_unlock(ser);
986
987 ast_threadstorage_set_ptr(&current_taskpool_serializer, prior_serializer);
988
989 return sync_task.fail;
990}
static void * listener(void *unused)
Definition asterisk.c:1531
void * ao2_object_get_lockaddr(void *obj)
Return the mutex lock address of an object.
Definition astobj2.c:476
A listener for taskprocessors.
A ast_taskprocessor structure is a singleton by name.
ast_cond_t cond
Definition taskpool.c:725
enum serializer_suspension suspended
Definition taskpool.c:727
struct ast_taskprocessor * ast_taskpool_serializer_get_current(void)
Get the taskpool serializer currently associated with this thread.
Definition taskpool.c:842
static int taskpool_serializer_empty_task(void *data)
Definition taskpool.c:885
void * ast_taskprocessor_listener_get_user_data(const struct ast_taskprocessor_listener *listener)
Get the user data from the listener.
long ast_taskprocessor_size(struct ast_taskprocessor *tps)
Return the current size of the taskprocessor queue.
int ast_taskprocessor_execute(struct ast_taskprocessor *tps)
Pop a task off the taskprocessor and execute it.
#define ast_taskprocessor_push(tps, task_exe, datap)
int ast_threadstorage_set_ptr(struct ast_threadstorage *ts, void *ptr)
Set a raw pointer from threadstorage.

References __ast_taskprocessor_push(), ao2_lock, ao2_object_get_lockaddr(), ao2_unlock, ast_cond_wait, ast_mutex_lock, ast_mutex_unlock, ast_taskpool_get_current(), ast_taskpool_serializer_get_current(), ast_taskprocessor_execute(), ast_taskprocessor_listener_get_user_data(), ast_taskprocessor_push, ast_taskprocessor_size(), ast_threadstorage_set_ptr(), taskpool_sync_task::complete, taskpool_sync_task::cond, serializer::cond, taskpool_sync_task::fail, listener(), taskpool_sync_task::lock, NULL, SERIALIZER_SUSPENDED, serializer::suspended, task(), taskpool_serializer_empty_task(), taskpool_sync_task_cleanup(), and taskpool_sync_task_init().

Referenced by __ast_sip_push_task_wait(), __ast_sip_push_task_wait_serializer(), and ast_taskpool_serializer_push_wait().

◆ ast_taskpool_create()

struct ast_taskpool * ast_taskpool_create ( const char *  name,
const struct ast_taskpool_options options 
)

Create a new taskpool.

Since
23.1.0
22.7.0
20.17.0

This function creates a taskpool. Tasks may be pushed onto this task pool and will be automatically acted upon by taskprocessors within the pool.

Only a single taskpool with a given name may exist. This function will fail if a taskpool with the given name already exists.

Parameters
nameThe unique name for the taskpool
optionsThe behavioral options for this taskpool
Return values
NULLFailed to create the taskpool
non-NULLThe newly-created taskpool
Note
The ast_taskpool_shutdown function must be called to shut down the taskpool and clean up underlying resources fully.

Definition at line 341 of file taskpool.c.

343{
344 struct ast_taskpool *pool;
345
346 /* Enforce versioning on the passed-in options */
347 if (options->version != AST_TASKPOOL_OPTIONS_VERSION) {
348 return NULL;
349 }
350
351 pool = ao2_alloc(sizeof(*pool) + strlen(name) + 1, NULL);
352 if (!pool) {
353 return NULL;
354 }
355
356 strcpy(pool->name, name); /* Safe */
357 memcpy(&pool->options, options, sizeof(pool->options));
358 pool->shrink_sched_id = -1;
359
360 /* Verify the passed-in options are valid, and adjust if needed */
361 if (options->initial_size < options->minimum_size) {
362 pool->options.initial_size = options->minimum_size;
363 ast_log(LOG_WARNING, "Taskpool '%s' has an initial size of %d, which is less than the minimum size of %d. Adjusting to %d.\n",
364 name, options->initial_size, options->minimum_size, options->minimum_size);
365 }
366
367 if (options->max_size && pool->options.initial_size > options->max_size) {
368 pool->options.max_size = pool->options.initial_size;
369 ast_log(LOG_WARNING, "Taskpool '%s' has a max size of %d, which is less than the initial size of %d. Adjusting to %d.\n",
370 name, options->max_size, pool->options.initial_size, pool->options.initial_size);
371 }
372
373 if (!options->auto_increment) {
374 if (!pool->options.minimum_size) {
375 pool->options.minimum_size = 1;
376 ast_log(LOG_WARNING, "Taskpool '%s' has a minimum size of 0, which is not valid without auto increment. Adjusting to 1.\n", name);
377 }
378 if (!pool->options.max_size) {
379 pool->options.max_size = pool->options.minimum_size;
380 ast_log(LOG_WARNING, "Taskpool '%s' has a max size of 0, which is not valid without auto increment. Adjusting to %d.\n", name, pool->options.minimum_size);
381 }
382 if (pool->options.minimum_size != pool->options.max_size) {
383 pool->options.minimum_size = pool->options.max_size;
384 pool->options.initial_size = pool->options.max_size;
385 ast_log(LOG_WARNING, "Taskpool '%s' has a minimum size of %d, while max size is %d. Adjusting all sizes to %d due to lack of auto increment.\n",
386 name, options->minimum_size, pool->options.max_size, pool->options.max_size);
387 }
388 } else if (!options->growth_threshold) {
390 }
391
394 } else if (options->selector == AST_TASKPOOL_SELECTOR_SEQUENTIAL) {
396 } else {
397 ast_log(LOG_WARNING, "Taskpool '%s' has an invalid selector of %d. Adjusting to default selector.\n",
398 name, options->selector);
400 }
401
403 ao2_ref(pool, -1);
404 return NULL;
405 }
406
407 /* Create the static taskprocessors based on the passed-in options */
408 for (int i = 0; i < pool->options.minimum_size; i++) {
410
412 if (!taskprocessor) {
413 /* The reference to pool is passed to ast_taskpool_shutdown */
415 return NULL;
416 }
417
420 /* The reference to pool is passed to ast_taskpool_shutdown */
422 return NULL;
423 }
424 }
425
427 pool->options.initial_size - pool->options.minimum_size)) {
429 return NULL;
430 }
431
432 /* Create the dynamic taskprocessor based on the passed-in options */
433 for (int i = 0; i < (pool->options.initial_size - pool->options.minimum_size); i++) {
435
437 if (!taskprocessor) {
438 /* The reference to pool is passed to ast_taskpool_shutdown */
440 return NULL;
441 }
442
445 /* The reference to pool is passed to ast_taskpool_shutdown */
447 return NULL;
448 }
449 }
450
451 /* If idle timeout support is enabled kick off a scheduled task to shrink the dynamic pool periodically, we do
452 * this no matter if there are dynamic taskprocessor present to reduce the work needed within the push function
453 * and to reduce complexity.
454 */
455 if (options->idle_timeout && options->auto_increment) {
457 if (pool->shrink_sched_id < 0) {
458 ao2_ref(pool, -1);
459 /* The second reference to pool is passed to ast_taskpool_shutdown */
461 return NULL;
462 }
463 }
464
465 return pool;
466}
#define ast_log
Definition astobj2.c:42
#define ao2_ref(o, delta)
Reference/unreference an object and return the old refcount.
Definition astobj2.h:459
#define ao2_alloc(data_size, destructor_fn)
Definition astobj2.h:409
static const char name[]
Definition format_mp3.c:68
#define LOG_WARNING
int ast_sched_add(struct ast_sched_context *con, int when, ast_sched_cb callback, const void *data) attribute_warn_unused_result
Adds a scheduled event.
Definition sched.c:567
int max_size
Maximum number of taskprocessors a pool may have.
Definition taskpool.h:122
int growth_threshold
The threshold for when to grow the pool.
Definition taskpool.h:131
int minimum_size
Number of taskprocessors that will always exist.
Definition taskpool.h:99
int initial_size
Number of taskprocessors the pool will start with.
Definition taskpool.h:109
An opaque taskpool structure.
Definition taskpool.c:62
int shrink_sched_id
Definition taskpool.c:72
char name[0]
Definition taskpool.c:76
Definition sched.c:76
struct ast_taskprocessor * taskprocessor
Definition taskpool.c:36
struct taskpool_taskprocessors::@436 taskprocessors
static void taskpool_least_full_selector(struct ast_taskpool *pool, struct taskpool_taskprocessors *taskprocessors, struct taskpool_taskprocessor **taskprocessor, unsigned int *growth_threshold_reached)
Least full taskprocessor selector.
Definition taskpool.c:301
void ast_taskpool_shutdown(struct ast_taskpool *pool)
Shut down a taskpool and remove the underlying taskprocessors.
Definition taskpool.c:692
static int taskpool_dynamic_pool_shrink(const void *data)
Definition taskpool.c:240
static int taskpool_taskprocessors_init(struct taskpool_taskprocessors *taskprocessors, unsigned int size)
Definition taskpool.c:207
static struct taskpool_taskprocessor * taskpool_taskprocessor_alloc(struct ast_taskpool *pool, char type)
Definition taskpool.c:166
static void taskpool_sequential_selector(struct ast_taskpool *pool, struct taskpool_taskprocessors *taskprocessors, struct taskpool_taskprocessor **taskprocessor, unsigned int *growth_threshold_reached)
Definition taskpool.c:276
#define TASKPOOL_GROW_THRESHOLD
The threshold for a taskprocessor at which we consider the pool needing to grow (50% of high water th...
Definition taskpool.c:80
#define AST_TASKPOOL_OPTIONS_VERSION
Definition taskpool.h:76
@ AST_TASKPOOL_SELECTOR_SEQUENTIAL
Definition taskpool.h:72
@ AST_TASKPOOL_SELECTOR_LEAST_FULL
Definition taskpool.h:71
@ AST_TASKPOOL_SELECTOR_DEFAULT
Definition taskpool.h:70
static struct test_options options
#define AST_VECTOR_APPEND(vec, elem)
Append an element to a vector, growing the vector if needed.
Definition vector.h:267

References ao2_alloc, ao2_bump, ao2_ref, ast_log, ast_sched_add(), AST_TASKPOOL_OPTIONS_VERSION, AST_TASKPOOL_SELECTOR_DEFAULT, AST_TASKPOOL_SELECTOR_LEAST_FULL, AST_TASKPOOL_SELECTOR_SEQUENTIAL, ast_taskpool_shutdown(), AST_VECTOR_APPEND, ast_taskpool::dynamic_taskprocessors, ast_taskpool_options::growth_threshold, ast_taskpool_options::initial_size, LOG_WARNING, ast_taskpool_options::max_size, ast_taskpool_options::minimum_size, ast_taskpool::name, name, NULL, ast_taskpool::options, options, ast_taskpool::selector, ast_taskpool::shrink_sched_id, ast_taskpool::static_taskprocessors, taskpool_dynamic_pool_shrink(), TASKPOOL_GROW_THRESHOLD, taskpool_least_full_selector(), taskpool_sequential_selector(), taskpool_taskprocessor_alloc(), taskpool_taskprocessors_init(), taskpool_taskprocessor::taskprocessor, and taskpool_taskprocessors::taskprocessors.

Referenced by ast_sorcery_init(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), handle_cli_taskpool_push_efficiency(), handle_cli_taskpool_push_serializer_efficiency(), load_module(), sel_run_case(), and stasis_init().

◆ ast_taskpool_get_current()

static struct ast_taskpool * ast_taskpool_get_current ( void  )
static

Definition at line 106 of file taskpool.c.

107{
108 return ast_threadstorage_get_ptr(&current_taskpool_pool);
109}
void * ast_threadstorage_get_ptr(struct ast_threadstorage *ts)
Retrieve a raw pointer from threadstorage.

References ast_threadstorage_get_ptr().

Referenced by __ast_taskpool_push_wait(), __ast_taskpool_serializer_push_wait(), ast_taskpool_serializer_suspend(), ast_taskpool_serializer_unsuspend(), execute_tasks(), and taskpool_taskprocessor_stop().

◆ ast_taskpool_init()

int ast_taskpool_init ( void  )

Provided by taskpool.c

Definition at line 1110 of file taskpool.c.

1111{
1113 if (!sched) {
1114 return -1;
1115 }
1116
1118 return -1;
1119 }
1120
1122
1123 return 0;
1124}
int ast_register_cleanup(void(*func)(void))
Register a function to be executed before Asterisk gracefully exits.
Definition clicompat.c:19
int ast_sched_start_thread(struct ast_sched_context *con)
Start a thread for processing scheduler entries.
Definition sched.c:197
struct ast_sched_context * ast_sched_context_create(void)
Create a scheduler context.
Definition sched.c:238
static void taskpool_shutdown(void)
Definition taskpool.c:1102

References ast_register_cleanup(), ast_sched_context_create(), ast_sched_start_thread(), and taskpool_shutdown().

Referenced by asterisk_daemon().

◆ ast_taskpool_push()

int ast_taskpool_push ( struct ast_taskpool pool,
int(*)(void *data)  task,
void *  data 
)

Definition at line 589 of file taskpool.c.

590{
591 return __ast_taskpool_push(pool, task, data, NULL, 0, NULL);
592}

References __ast_taskpool_push(), NULL, and task().

◆ ast_taskpool_push_wait()

int ast_taskpool_push_wait ( struct ast_taskpool pool,
int(*)(void *data)  task,
void *  data 
)

Definition at line 687 of file taskpool.c.

688{
689 return __ast_taskpool_push_wait(pool, task, data, NULL, 0, NULL);
690}
int __ast_taskpool_push_wait(struct ast_taskpool *pool, int(*task)(void *data), void *data, const char *file, int line, const char *function)
Push a task to the taskpool, and wait for completion.
Definition taskpool.c:652

References __ast_taskpool_push_wait(), NULL, and task().

◆ ast_taskpool_queue_size()

long ast_taskpool_queue_size ( struct ast_taskpool pool)

Get the current number of queued tasks in the taskpool.

Since
23.1.0
22.7.0
20.17.0
Parameters
poolThe taskpool to query
Return values
Thenumber of queued tasks in the taskpool

Definition at line 481 of file taskpool.c.

482{
483 long queue_size = 0;
484
485 ao2_lock(pool);
488 ao2_unlock(pool);
489
490 return queue_size;
491}
#define TASKPOOL_QUEUE_SIZE_ADD(tps, size)
Definition taskpool.c:479
#define AST_VECTOR_CALLBACK_VOID(vec, callback,...)
Execute a callback on every element in a vector disregarding callback return.
Definition vector.h:890

References ao2_lock, ao2_unlock, AST_VECTOR_CALLBACK_VOID, ast_taskpool::dynamic_taskprocessors, ast_taskpool::static_taskprocessors, TASKPOOL_QUEUE_SIZE_ADD, and taskpool_taskprocessors::taskprocessors.

Referenced by ast_sip_taskpool_queue_size().

◆ ast_taskpool_serializer()

struct ast_taskprocessor * ast_taskpool_serializer ( const char *  name,
struct ast_taskpool pool 
)

Serialized execution of tasks within a ast_taskpool.

Since
23.1.0
22.7.0
20.17.0

A ast_taskprocessor with the same contract as a default taskprocessor (tasks execute serially) except instead of executing out of a dedicated thread, execution occurs in a taskprocessor from a ast_taskpool.

While it guarantees that each task will complete before executing the next, there is no guarantee as to which thread from the pool individual tasks will execute. This normally only matters if your code relies on thread specific information, such as thread locals.

Use ast_taskprocessor_unreference() to dispose of the returned ast_taskprocessor.

Only a single taskprocessor with a given name may exist. This function will fail if a taskprocessor with the given name already exists.

Parameters
nameName of the serializer. (must be unique)
poolast_taskpool for execution.
Returns
ast_taskprocessor for enqueuing work.
Return values
NULLon error.

Definition at line 877 of file taskpool.c.

878{
880}
struct ast_taskprocessor * ast_taskpool_serializer_group(const char *name, struct ast_taskpool *pool, struct ast_serializer_shutdown_group *shutdown_group)
Serialized execution of tasks within a ast_taskpool.
Definition taskpool.c:847

References ast_taskpool_serializer_group(), name, and NULL.

Referenced by AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), handle_cli_taskpool_push_serializer_efficiency(), internal_stasis_subscribe(), sel_run_case(), and sorcery_object_type_alloc().

◆ ast_taskpool_serializer_get_current()

struct ast_taskprocessor * ast_taskpool_serializer_get_current ( void  )

Get the taskpool serializer currently associated with this thread.

Note
The returned pointer is valid while the serializer thread is running.
Use ao2_ref() on serializer if you are going to keep it for another thread. To unref it you must then use ast_taskprocessor_unreference().
Return values
serializeron success.
NULLon error or no serializer associated with the thread.

Definition at line 842 of file taskpool.c.

843{
844 return ast_threadstorage_get_ptr(&current_taskpool_serializer);
845}

References ast_threadstorage_get_ptr().

Referenced by __ast_taskpool_serializer_push_wait(), record_serializer(), requeue_task(), rfc3326_outgoing_request(), rfc3326_outgoing_response(), serializer_efficiency_task(), and simple_task().

◆ ast_taskpool_serializer_group()

struct ast_taskprocessor * ast_taskpool_serializer_group ( const char *  name,
struct ast_taskpool pool,
struct ast_serializer_shutdown_group shutdown_group 
)

Serialized execution of tasks within a ast_taskpool.

Since
23.1.0
22.7.0
20.17.0

A ast_taskprocessor with the same contract as a default taskprocessor (tasks execute serially) except instead of executing out of a dedicated thread, execution occurs in a taskprocessor from a ast_taskpool.

While it guarantees that each task will complete before executing the next, there is no guarantee as to which thread from the pool individual tasks will execute. This normally only matters if your code relies on thread specific information, such as thread locals.

Use ast_taskprocessor_unreference() to dispose of the returned ast_taskprocessor.

Only a single taskprocessor with a given name may exist. This function will fail if a taskprocessor with the given name already exists.

Parameters
nameName of the serializer. (must be unique)
poolast_taskpool for execution.
shutdown_groupGroup shutdown controller. (NULL if no group association)
Returns
ast_taskprocessor for enqueuing work.
Return values
NULLon error.

Definition at line 847 of file taskpool.c.

849{
850 struct serializer *ser;
852 struct ast_taskprocessor *tps;
853
855 if (!ser) {
856 return NULL;
857 }
858
860 if (!listener) {
861 ao2_ref(ser, -1);
862 return NULL;
863 }
864
866 if (!tps) {
867 /* ser ref transferred to listener but not cleaned without tps */
868 ao2_ref(ser, -1);
869 } else if (shutdown_group) {
871 }
872
873 ao2_ref(listener, -1);
874 return tps;
875}
static struct ast_serializer_shutdown_group * shutdown_group
void ast_serializer_shutdown_group_inc(struct ast_serializer_shutdown_group *shutdown_group)
Increment the number of serializer members in the group.
static struct ast_taskprocessor_listener_callbacks serializer_tps_listener_callbacks
Definition taskpool.c:836
static struct serializer * serializer_create(struct ast_taskpool *pool, struct ast_serializer_shutdown_group *shutdown_group)
Definition taskpool.c:739
struct ast_taskprocessor_listener * ast_taskprocessor_listener_alloc(const struct ast_taskprocessor_listener_callbacks *callbacks, void *user_data)
Allocate a taskprocessor listener.
struct ast_taskprocessor * ast_taskprocessor_create_with_listener(const char *name, struct ast_taskprocessor_listener *listener)
Create a taskprocessor with a custom listener.

References ao2_ref, ast_serializer_shutdown_group_inc(), ast_taskprocessor_create_with_listener(), ast_taskprocessor_listener_alloc(), listener(), name, NULL, serializer_create(), serializer_tps_listener_callbacks, and shutdown_group.

Referenced by ast_serializer_taskpool_create(), ast_sip_create_serializer_group(), and ast_taskpool_serializer().

◆ ast_taskpool_serializer_push_wait()

int ast_taskpool_serializer_push_wait ( struct ast_taskprocessor serializer,
int(*)(void *data)  task,
void *  data 
)

Definition at line 894 of file taskpool.c.

895{
897}
int __ast_taskpool_serializer_push_wait(struct ast_taskprocessor *serializer, int(*task)(void *data), void *data, const char *file, int line, const char *function)
Push a task to a serializer, and wait for completion.
Definition taskpool.c:899

References __ast_taskpool_serializer_push_wait(), NULL, and task().

◆ ast_taskpool_serializer_suspend()

int ast_taskpool_serializer_suspend ( struct ast_taskprocessor serializer)

Suspend a serializer, causing tasks to be queued until unsuspended.

Since
23.2.0
22.8.0
20.18.0
Parameters
serializerThe serializer to suspend
Return values
0success
-1failure
Note
May only be invoked from outside of the taskpool

Definition at line 1018 of file taskpool.c.

1019{
1022
1023 /* This suspension process works by inserting a checkpoint into the queue of the
1024 * serializer. Once this checkpoint is reached the taskpool taskprocessor handling
1025 * the queue stops prematurely and does not get requeued. For the case where a
1026 * synchronous task wait is in progress it is instead paused temporarily. Once
1027 * the serializer is unsuspended a new execution task is queued into the taskpool
1028 * to resume execution and any paused synchronous task waits are awoken to resume
1029 * their own execution as well. This approach minimizes the number of threads that
1030 * are paused waiting, nominally to 0.
1031 */
1032
1034 return -1;
1035 }
1036
1037 ao2_lock(ser);
1038
1039 /* If the serializer is already suspending or suspended, just return immediately.
1040 * This mirrors the original behavior from PJSIP.
1041 */
1042 if (ser->suspended != SERIALIZER_UNSUSPENDED) {
1043 ao2_unlock(ser);
1044 return 0;
1045 }
1046
1048
1049 ao2_unlock(ser);
1050
1051 /* Once this returns successfully there is no thread executing the tasks on the serializer,
1052 * so they will accumulate until the serializer is unsuspended.
1053 */
1055 /* Suspension failed, so unsuspend as doing otherwise would leave the serializer in a stuck
1056 * state.
1057 */
1058 ao2_lock(ser);
1060 ao2_unlock(ser);
1061 return -1;
1062 }
1063
1064 return 0;
1065}
static int taskpool_serializer_suspend_task(void *data)
Definition taskpool.c:995
#define ast_taskpool_serializer_push_wait(pool, task, data)
Definition taskpool.h:335

References ao2_lock, ao2_unlock, ast_taskpool_get_current(), ast_taskpool_serializer_push_wait, ast_taskprocessor_listener_get_user_data(), listener(), SERIALIZER_SUSPENDING, SERIALIZER_UNSUSPENDED, serializer::suspended, and taskpool_serializer_suspend_task().

Referenced by ast_sip_session_suspend(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), and AST_TEST_DEFINE().

◆ ast_taskpool_serializer_unsuspend()

int ast_taskpool_serializer_unsuspend ( struct ast_taskprocessor serializer)

Unsuspend a serializer, causing tasks to be executed.

Since
23.2.0
22.8.0
20.18.0
Parameters
serializerThe serializer to unsuspend
Return values
0success
-1failure
Note
May only be invoked from outside of the taskpool

Definition at line 1067 of file taskpool.c.

1068{
1071
1073 return -1;
1074 }
1075
1076 ao2_lock(ser);
1077
1078 if (ser->suspended != SERIALIZER_SUSPENDED) {
1079 ao2_unlock(ser);
1080 return 0;
1081 }
1082
1084
1085 /* Notify any other interested threads that this one has awoken */
1086 ast_cond_broadcast(&ser->cond);
1087
1088 /* And now we kick off handling of the queued tasks once again */
1091 }
1092
1093 ao2_unlock(ser);
1094
1095 return 0;
1096}
#define ast_cond_broadcast(cond)
Definition lock.h:211
struct ast_taskpool * pool
Definition taskpool.c:721
static int execute_tasks(void *data)
Definition taskpool.c:759
#define ast_taskpool_push(pool, task, data)
Definition taskpool.h:210
void * ast_taskprocessor_unreference(struct ast_taskprocessor *tps)
Unreference the specified taskprocessor and its reference count will decrement.

References ao2_bump, ao2_lock, ao2_unlock, ast_cond_broadcast, ast_taskpool_get_current(), ast_taskpool_push, ast_taskprocessor_listener_get_user_data(), ast_taskprocessor_unreference(), serializer::cond, execute_tasks(), listener(), serializer::pool, SERIALIZER_SUSPENDED, SERIALIZER_UNSUSPENDED, and serializer::suspended.

Referenced by ast_sip_session_unsuspend(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), and AST_TEST_DEFINE().

◆ ast_taskpool_shutdown()

void ast_taskpool_shutdown ( struct ast_taskpool pool)

Shut down a taskpool and remove the underlying taskprocessors.

Since
23.1.0
22.7.0
20.17.0
Parameters
poolThe pool to shut down
Note
This will decrement the reference to the pool

Definition at line 692 of file taskpool.c.

693{
694 if (!pool) {
695 return;
696 }
697
698 /* Mark this pool as shutting down so nothing new is pushed */
699 ao2_lock(pool);
700 pool->shutting_down = 1;
701 ao2_unlock(pool);
702
703 /* Stop the shrink scheduled item if present */
705
706 /* Clean up all the taskprocessors */
709
710 ao2_ref(pool, -1);
711}
#define AST_SCHED_DEL_UNREF(sched, id, refcall)
schedule task to get deleted and call unref function
Definition sched.h:82
static void taskpool_taskprocessors_cleanup(struct taskpool_taskprocessors *taskprocessors)
Definition taskpool.c:220

References ao2_lock, ao2_ref, ao2_unlock, AST_SCHED_DEL_UNREF, ast_taskpool::dynamic_taskprocessors, ast_taskpool::shrink_sched_id, ast_taskpool::shutting_down, ast_taskpool::static_taskprocessors, and taskpool_taskprocessors_cleanup().

Referenced by ast_taskpool_create(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), AST_TEST_DEFINE(), handle_cli_taskpool_push_efficiency(), handle_cli_taskpool_push_serializer_efficiency(), load_module(), sel_run_case(), sorcery_cleanup(), stasis_cleanup(), and unload_module().

◆ ast_taskpool_taskprocessors_count()

size_t ast_taskpool_taskprocessors_count ( struct ast_taskpool pool)

Get the current number of taskprocessors in the taskpool.

Since
23.1.0
22.7.0
20.17.0
Parameters
poolThe taskpool to query
Return values
Thenumber of taskprocessors in the taskpool

Definition at line 468 of file taskpool.c.

469{
470 size_t count;
471
472 ao2_lock(pool);
474 ao2_unlock(pool);
475
476 return count;
477}
#define AST_VECTOR_SIZE(vec)
Get the number of elements in a vector.
Definition vector.h:637

References ao2_lock, ao2_unlock, AST_VECTOR_SIZE, ast_taskpool::dynamic_taskprocessors, ast_taskpool::static_taskprocessors, and taskpool_taskprocessors::taskprocessors.

Referenced by AST_TEST_DEFINE(), and AST_TEST_DEFINE().

◆ AST_THREADSTORAGE_RAW() [1/2]

AST_THREADSTORAGE_RAW ( current_taskpool_pool  )

Thread storage for the current taskpool.

◆ AST_THREADSTORAGE_RAW() [2/2]

AST_THREADSTORAGE_RAW ( current_taskpool_serializer  )

◆ execute_tasks()

static int execute_tasks ( void *  data)
static

Definition at line 759 of file taskpool.c.

760{
761 struct ast_taskpool *pool = ast_taskpool_get_current();
762 struct ast_taskprocessor *tps = data;
765 size_t remaining, requeue = 0;
766
767 /* In a normal scenario this lock will not be in contention with
768 * anything else. It is only if a synchronous task is pushed to
769 * the serializer that it may be blocked on the synchronous
770 * task thread. This is done to ensure that only one thread is executing
771 * tasks from the serializer at a given time, and not out of order
772 * either.
773 */
774 ao2_lock(ser);
775
776 ast_threadstorage_set_ptr(&current_taskpool_serializer, tps);
777 for (remaining = ast_taskprocessor_size(tps); remaining > 0; remaining--) {
778 requeue = ast_taskprocessor_execute(tps);
779
780 /* If the serializer is suspended we will not execute any more tasks and
781 * we will not requeue the taskpool task. Instead it will be requeued when
782 * the serializer is unsuspended.
783 */
784 if (ser->suspended == SERIALIZER_SUSPENDED) {
785 requeue = 0;
786 break;
787 }
788 }
789 ast_threadstorage_set_ptr(&current_taskpool_serializer, NULL);
790
791 ao2_unlock(ser);
792
793 /* If there are remaining tasks we requeue, this way the serializer
794 * does not hold exclusivity of the taskpool taskprocessor
795 */
796 if (requeue) {
797 /* Ownership passes to the new task */
800 }
801 } else {
803 }
804
805 return 0;
806}
struct ast_taskprocessor * tps

References ao2_lock, ao2_unlock, ast_taskpool_get_current(), ast_taskpool_push, ast_taskprocessor_execute(), ast_taskprocessor_listener_get_user_data(), ast_taskprocessor_size(), ast_taskprocessor_unreference(), ast_threadstorage_set_ptr(), execute_tasks(), listener(), NULL, serializer::pool, SERIALIZER_SUSPENDED, serializer::suspended, and ast_taskprocessor_listener::tps.

Referenced by ast_taskpool_serializer_unsuspend(), execute_tasks(), and serializer_task_pushed().

◆ serializer_create()

static struct serializer * serializer_create ( struct ast_taskpool pool,
struct ast_serializer_shutdown_group shutdown_group 
)
static

Definition at line 739 of file taskpool.c.

741{
742 struct serializer *ser;
743
744 /* This object has a lock so it can be used to ensure exclusive access
745 * to the execution of tasks within the serializer.
746 */
747 ser = ao2_alloc(sizeof(*ser), serializer_dtor);
748 if (!ser) {
749 return NULL;
750 }
751 ser->pool = ao2_bump(pool);
753 ast_cond_init(&ser->cond, NULL);
754 return ser;
755}
#define ast_cond_init(cond, attr)
Definition lock.h:208
struct ast_serializer_shutdown_group * shutdown_group
Definition taskpool.c:723
static void serializer_dtor(void *obj)
Definition taskpool.c:730

References ao2_alloc, ao2_bump, ast_cond_init, serializer::cond, NULL, serializer::pool, serializer_dtor(), serializer::shutdown_group, and shutdown_group.

Referenced by ast_taskpool_serializer_group().

◆ serializer_dtor()

static void serializer_dtor ( void *  obj)
static

Definition at line 730 of file taskpool.c.

731{
732 struct serializer *ser = obj;
733
734 ao2_cleanup(ser->pool);
736 ast_cond_destroy(&ser->cond);
737}
#define ast_cond_destroy(cond)
Definition lock.h:209

References ao2_cleanup, ast_cond_destroy, serializer::cond, serializer::pool, and serializer::shutdown_group.

Referenced by serializer_create().

◆ serializer_shutdown()

static void serializer_shutdown ( struct ast_taskprocessor_listener listener)
static

Definition at line 826 of file taskpool.c.

827{
829
830 if (ser->shutdown_group) {
832 }
833 ao2_cleanup(ser);
834}
void ast_serializer_shutdown_group_dec(struct ast_serializer_shutdown_group *shutdown_group)
Decrement the number of serializer members in the group.

References ao2_cleanup, ast_serializer_shutdown_group_dec(), ast_taskprocessor_listener_get_user_data(), listener(), and serializer::shutdown_group.

◆ serializer_start()

static int serializer_start ( struct ast_taskprocessor_listener listener)
static

Definition at line 820 of file taskpool.c.

821{
822 /* No-op */
823 return 0;
824}

◆ serializer_task_pushed()

static void serializer_task_pushed ( struct ast_taskprocessor_listener listener,
int  was_empty 
)
static

Definition at line 808 of file taskpool.c.

809{
810 if (was_empty) {
813
814 if (ast_taskpool_push(ser->pool, execute_tasks, tps)) {
816 }
817 }
818}
struct ast_taskprocessor * ast_taskprocessor_listener_get_tps(const struct ast_taskprocessor_listener *listener)
Get a reference to the listener's taskprocessor.

References ast_taskpool_push, ast_taskprocessor_listener_get_tps(), ast_taskprocessor_listener_get_user_data(), ast_taskprocessor_unreference(), execute_tasks(), listener(), and serializer::pool.

◆ taskpool_dynamic_pool_grow()

static void taskpool_dynamic_pool_grow ( struct ast_taskpool pool,
struct taskpool_taskprocessor **  taskprocessor 
)
static

Definition at line 496 of file taskpool.c.

497{
498 unsigned int num_to_add = pool->options.auto_increment;
499 int i;
500
501 if (!num_to_add) {
502 return;
503 }
504
505 /* If a maximum size is enforced, then determine if we have to limit how many taskprocessors we add */
506 if (pool->options.max_size) {
508
509 if (current_size + num_to_add > pool->options.max_size) {
510 num_to_add = pool->options.max_size - current_size;
511 }
512 }
513
514 for (i = 0; i < num_to_add; i++) {
515 struct taskpool_taskprocessor *new_taskprocessor;
516
517 new_taskprocessor = taskpool_taskprocessor_alloc(pool, 'd');
518 if (!new_taskprocessor) {
519 return;
520 }
521
522 if (AST_VECTOR_APPEND(&pool->dynamic_taskprocessors.taskprocessors, new_taskprocessor)) {
523 ao2_ref(new_taskprocessor, -1);
524 return;
525 }
526
527 if (i == 0) {
528 /* On the first iteration we return the taskprocessor we just added */
529 *taskprocessor = new_taskprocessor;
530 /* We assume we will be going back to the first taskprocessor, since we are at the end of the vector */
532 } else if (i == 1) {
533 /* On the second iteration we update the next taskprocessor to use to be this one */
535 }
536 }
537}
unsigned int taskprocessor_num
Definition taskpool.c:48

References ao2_ref, AST_VECTOR_APPEND, AST_VECTOR_SIZE, ast_taskpool_options::auto_increment, ast_taskpool::dynamic_taskprocessors, ast_taskpool_options::max_size, ast_taskpool::options, ast_taskpool::static_taskprocessors, taskpool_taskprocessor_alloc(), taskpool_taskprocessor::taskprocessor, taskpool_taskprocessors::taskprocessor_num, and taskpool_taskprocessors::taskprocessors.

Referenced by __ast_taskpool_push().

◆ taskpool_dynamic_pool_shrink()

static int taskpool_dynamic_pool_shrink ( const void *  data)
static

Definition at line 240 of file taskpool.c.

241{
242 struct ast_taskpool *pool = (struct ast_taskpool *)data;
243 int num_removed;
244
245 ao2_lock(pool);
246
247 /* If the pool is shutting down, do nothing and don't reschedule */
248 if (pool->shutting_down) {
249 ao2_unlock(pool);
250 ao2_ref(pool, -1);
251 return 0;
252 }
253
254 /* Go through the dynamic taskprocessors and find any which have been idle long enough and remove them */
257 if (num_removed) {
258 /* If we've removed any taskprocessors the taskprocessor_num may no longer be valid, so update it */
261 }
262 }
263
264 ao2_unlock(pool);
265
266 /* It is possible for the pool to have been shut down between unlocking and returning, this is
267 * inherently a race condition we can't eliminate so we will catch it on the next iteration.
268 */
269 return pool->options.idle_timeout * 1000;
270}
int idle_timeout
Time limit in seconds for idle dynamic taskprocessors.
Definition taskpool.h:88
#define TASKPROCESSOR_IS_IDLE(tps, timeout)
Definition taskpool.c:235
#define AST_VECTOR_REMOVE_ALL_CMP_UNORDERED(vec, value, cmp, cleanup)
Remove all elements from a vector that matches the given comparison.
Definition vector.h:489

References ao2_cleanup, ao2_lock, ao2_ref, ao2_unlock, AST_VECTOR_REMOVE_ALL_CMP_UNORDERED, AST_VECTOR_SIZE, ast_taskpool::dynamic_taskprocessors, ast_taskpool_options::idle_timeout, ast_taskpool::options, ast_taskpool::shutting_down, TASKPROCESSOR_IS_IDLE, taskpool_taskprocessors::taskprocessor_num, and taskpool_taskprocessors::taskprocessors.

Referenced by ast_taskpool_create().

◆ taskpool_least_full_selector()

static void taskpool_least_full_selector ( struct ast_taskpool pool,
struct taskpool_taskprocessors taskprocessors,
struct taskpool_taskprocessor **  taskprocessor,
unsigned int *  growth_threshold_reached 
)
static

Least full taskprocessor selector.

\interal

Definition at line 301 of file taskpool.c.

303{
304 struct taskpool_taskprocessor *least_full = NULL;
305 long least_full_load = 0;
306 unsigned int i;
307
308 if (!AST_VECTOR_SIZE(&taskprocessors->taskprocessors)) {
309 *growth_threshold_reached = 1;
310 return;
311 }
312
313 /* We assume that the growth threshold has not yet been reached, until proven otherwise */
314 *growth_threshold_reached = 0;
315
316 for (i = 0; i < AST_VECTOR_SIZE(&taskprocessors->taskprocessors); i++) {
317 struct taskpool_taskprocessor *tp = AST_VECTOR_GET(&taskprocessors->taskprocessors, i);
318 long load = taskpool_taskprocessor_load(tp);
319
320 /* If this taskprocessor has nothing queued and nothing in flight, it is the best choice */
321 if (!load) {
322 *taskprocessor = tp;
323 return;
324 }
325
326 /* If any of the taskprocessors have reached the growth threshold then we should grow the pool */
327 if (load >= pool->options.growth_threshold) {
328 *growth_threshold_reached = 1;
329 }
330
331 /* The taskprocessor with the lowest load should be used */
332 if (!least_full || load < least_full_load) {
333 least_full = tp;
334 least_full_load = load;
335 }
336 }
337
338 *taskprocessor = least_full;
339}
static long taskpool_taskprocessor_load(struct taskpool_taskprocessor *tp)
Definition taskpool.c:90
#define AST_VECTOR_GET(vec, idx)
Get an element from a vector.
Definition vector.h:708

References AST_VECTOR_GET, AST_VECTOR_SIZE, ast_taskpool_options::growth_threshold, NULL, ast_taskpool::options, taskpool_taskprocessor_load(), taskpool_taskprocessor::taskprocessor, and taskpool_taskprocessors::taskprocessors.

Referenced by ast_taskpool_create().

◆ taskpool_sequential_selector()

static void taskpool_sequential_selector ( struct ast_taskpool pool,
struct taskpool_taskprocessors taskprocessors,
struct taskpool_taskprocessor **  taskprocessor,
unsigned int *  growth_threshold_reached 
)
static

Definition at line 276 of file taskpool.c.

278{
279 unsigned int taskprocessor_num = taskprocessors->taskprocessor_num;
280
281 if (!AST_VECTOR_SIZE(&taskprocessors->taskprocessors)) {
282 *growth_threshold_reached = 1;
283 return;
284 }
285
286 taskprocessors->taskprocessor_num++;
287 if (taskprocessors->taskprocessor_num == AST_VECTOR_SIZE(&taskprocessors->taskprocessors)) {
288 taskprocessors->taskprocessor_num = 0;
289 }
290
291 *taskprocessor = AST_VECTOR_GET(&taskprocessors->taskprocessors, taskprocessor_num);
292
293 /* Check to see if this has reached the growth threshold */
294 *growth_threshold_reached = (taskpool_taskprocessor_load(*taskprocessor) >= pool->options.growth_threshold) ? 1 : 0;
295}

References AST_VECTOR_GET, AST_VECTOR_SIZE, ast_taskpool_options::growth_threshold, ast_taskpool::options, taskpool_taskprocessor_load(), taskpool_taskprocessors::taskprocessor_num, and taskpool_taskprocessors::taskprocessors.

Referenced by ast_taskpool_create().

◆ taskpool_serializer_empty_task()

static int taskpool_serializer_empty_task ( void *  data)
static

Definition at line 885 of file taskpool.c.

886{
887 return 0;
888}

Referenced by __ast_taskpool_serializer_push_wait(), and taskpool_serializer_suspend_task().

◆ taskpool_serializer_suspend_task()

static int taskpool_serializer_suspend_task ( void *  data)
static

Definition at line 995 of file taskpool.c.

996{
997 struct ast_taskprocessor *serializer = data;
1000
1001 /* First we queue the empty task to ensure the serializer doesn't reach empty, this
1002 * prevents any threads from queueing up a taskpool task that executes the serializer
1003 * while it is suspended, allowing us to queue it ourselves when the serializer is
1004 * unsuspended.
1005 */
1007 return -1;
1008 }
1009
1010 /* Next we suspend the serializer so that the execute_tasks currently executing stops
1011 * and doesn't requeue.
1012 */
1014
1015 return 0;
1016}

References ast_taskprocessor_listener_get_user_data(), ast_taskprocessor_push, listener(), NULL, SERIALIZER_SUSPENDED, serializer::suspended, and taskpool_serializer_empty_task().

Referenced by ast_taskpool_serializer_suspend().

◆ taskpool_shutdown()

static void taskpool_shutdown ( void  )
static

Definition at line 1102 of file taskpool.c.

1103{
1104 if (sched) {
1106 sched = NULL;
1107 }
1108}
void ast_sched_context_destroy(struct ast_sched_context *c)
destroys a schedule context
Definition sched.c:271

References ast_sched_context_destroy(), and NULL.

Referenced by ast_taskpool_init().

◆ taskpool_sync_task()

static int taskpool_sync_task ( void *  data)
static

Definition at line 632 of file taskpool.c.

633{
634 struct taskpool_sync_task *sync_task = data;
635 int ret;
636
637 sync_task->fail = sync_task->task(sync_task->task_data);
638
639 /*
640 * Once we unlock sync_task->lock after signaling, we cannot access
641 * sync_task again. The thread waiting within ast_taskpool_push_wait()
642 * is free to continue and release its local variable (sync_task).
643 */
644 ast_mutex_lock(&sync_task->lock);
645 sync_task->complete = 1;
646 ast_cond_signal(&sync_task->cond);
647 ret = sync_task->fail;
648 ast_mutex_unlock(&sync_task->lock);
649 return ret;
650}
#define ast_cond_signal(cond)
Definition lock.h:210
ast_cond_t cond
Definition taskpool.c:599
int(* task)(void *)
Definition taskpool.c:602
ast_mutex_t lock
Definition taskpool.c:598

References ast_cond_signal, ast_mutex_lock, ast_mutex_unlock, taskpool_sync_task::complete, taskpool_sync_task::cond, taskpool_sync_task::fail, taskpool_sync_task::lock, taskpool_sync_task::task, and taskpool_sync_task::task_data.

◆ taskpool_sync_task_cleanup()

static void taskpool_sync_task_cleanup ( struct taskpool_sync_task sync_task)
static

Definition at line 623 of file taskpool.c.

624{
625 ast_mutex_destroy(&sync_task->lock);
626 ast_cond_destroy(&sync_task->cond);
627}
#define ast_mutex_destroy(a)
Definition lock.h:195

References ast_cond_destroy, ast_mutex_destroy, taskpool_sync_task::cond, and taskpool_sync_task::lock.

Referenced by __ast_taskpool_push_wait(), and __ast_taskpool_serializer_push_wait().

◆ taskpool_sync_task_init()

static int taskpool_sync_task_init ( struct taskpool_sync_task sync_task,
int(*)(void *)  task,
void *  data 
)
static

Definition at line 609 of file taskpool.c.

610{
611 ast_mutex_init(&sync_task->lock);
612 ast_cond_init(&sync_task->cond, NULL);
613 sync_task->complete = 0;
614 sync_task->fail = 0;
615 sync_task->task = task;
616 sync_task->task_data = data;
617 return 0;
618}
#define ast_mutex_init(pmutex)
Definition lock.h:193

References ast_cond_init, ast_mutex_init, taskpool_sync_task::complete, taskpool_sync_task::cond, taskpool_sync_task::fail, taskpool_sync_task::lock, NULL, task(), taskpool_sync_task::task, and taskpool_sync_task::task_data.

Referenced by __ast_taskpool_push_wait(), and __ast_taskpool_serializer_push_wait().

◆ taskpool_taskprocessor_alloc()

static struct taskpool_taskprocessor * taskpool_taskprocessor_alloc ( struct ast_taskpool pool,
char  type 
)
static

Definition at line 166 of file taskpool.c.

167{
169 char tps_name[AST_TASKPROCESSOR_MAX_NAME + 1];
170
171 /* We don't actually need locking for each pool taskprocessor, as the only thing
172 * mutable is the underlying taskprocessor which has its own internal locking.
173 */
175 if (!taskprocessor) {
176 return NULL;
177 }
178
179 /* Create name with seq number appended. */
180 ast_taskprocessor_build_name(tps_name, sizeof(tps_name), "taskpool/%c:%s", type, pool->name);
181
182 taskprocessor->taskprocessor = ast_taskprocessor_get(tps_name, TPS_REF_DEFAULT);
183 if (!taskprocessor->taskprocessor) {
185 return NULL;
186 }
187
188 taskprocessor->last_pushed = ast_tvnow();
189
191 ao2_ref(pool, -1);
192 /* Prevent the taskprocessor from queueing the stop task by explicitly unreferencing and setting it to
193 * NULL here.
194 */
196 taskprocessor->taskprocessor = NULL;
197 return NULL;
198 }
199
200 return taskprocessor;
201}
@ AO2_ALLOC_OPT_LOCK_NOLOCK
Definition astobj2.h:367
#define ao2_alloc_options(data_size, destructor_fn, options)
Definition astobj2.h:404
static const char type[]
static int taskpool_taskprocessor_start(void *data)
Definition taskpool.c:145
static void taskpool_taskprocessor_dtor(void *obj)
Definition taskpool.c:130
struct ast_taskprocessor * ast_taskprocessor_get(const char *name, enum ast_tps_options create)
Get a reference to a taskprocessor with the specified name and create the taskprocessor if necessary.
@ TPS_REF_DEFAULT
return a reference to a taskprocessor, create one if it does not exist
void ast_taskprocessor_build_name(char *buf, unsigned int size, const char *format,...)
Build a taskprocessor name with a sequence number on the end.
#define AST_TASKPROCESSOR_MAX_NAME
Suggested maximum taskprocessor name length (less null terminator).

References AO2_ALLOC_OPT_LOCK_NOLOCK, ao2_alloc_options, ao2_bump, ao2_ref, ast_taskprocessor_build_name(), ast_taskprocessor_get(), AST_TASKPROCESSOR_MAX_NAME, ast_taskprocessor_push, ast_taskprocessor_unreference(), ast_tvnow(), ast_taskpool::name, NULL, taskpool_taskprocessor_dtor(), taskpool_taskprocessor_start(), taskpool_taskprocessor::taskprocessor, TPS_REF_DEFAULT, and type.

Referenced by ast_taskpool_create(), and taskpool_dynamic_pool_grow().

◆ taskpool_taskprocessor_dtor()

static void taskpool_taskprocessor_dtor ( void *  obj)
static

Definition at line 130 of file taskpool.c.

131{
133
135 /* We can't actually do anything if this fails, so just accept reality */
136 }
137
139}
static int taskpool_taskprocessor_stop(void *data)
Definition taskpool.c:115

References ast_taskprocessor_push, ast_taskprocessor_unreference(), NULL, taskpool_taskprocessor_stop(), and taskpool_taskprocessor::taskprocessor.

Referenced by taskpool_taskprocessor_alloc().

◆ taskpool_taskprocessor_load()

static long taskpool_taskprocessor_load ( struct taskpool_taskprocessor tp)
static

Definition at line 90 of file taskpool.c.

91{
94}
unsigned int ast_taskprocessor_is_executing(const struct ast_taskprocessor *tps)
Return whether the taskprocessor is currently executing a task.

References ast_taskprocessor_is_executing(), ast_taskprocessor_size(), and taskpool_taskprocessor::taskprocessor.

Referenced by taskpool_least_full_selector(), and taskpool_sequential_selector().

◆ taskpool_taskprocessor_start()

static int taskpool_taskprocessor_start ( void *  data)
static

Definition at line 145 of file taskpool.c.

146{
147 struct ast_taskpool *pool = data;
148
149 /* Set the pool on the thread for this taskprocessor, inheriting the
150 * reference passed to the task itself.
151 */
152 ast_threadstorage_set_ptr(&current_taskpool_pool, pool);
153
154 /* If a thread start callback is set on the options, call it */
155 if (pool->options.thread_start) {
156 pool->options.thread_start();
157 }
158
159 return 0;
160}
void(* thread_start)(void)
Function to call when a taskprocessor starts.
Definition taskpool.h:138

References ast_threadstorage_set_ptr(), ast_taskpool::options, and ast_taskpool_options::thread_start.

Referenced by taskpool_taskprocessor_alloc().

◆ taskpool_taskprocessor_stop()

static int taskpool_taskprocessor_stop ( void *  data)
static

Definition at line 115 of file taskpool.c.

116 {
117 struct ast_taskpool *pool = ast_taskpool_get_current();
118
119 /* If a thread stop callback is set on the options, call it */
120 if (pool->options.thread_end) {
121 pool->options.thread_end();
122 }
123
124 ao2_cleanup(pool);
125
126 return 0;
127 }
void(* thread_end)(void)
Function to call when a taskprocessor ends.
Definition taskpool.h:145

References ao2_cleanup, ast_taskpool_get_current(), ast_taskpool::options, and ast_taskpool_options::thread_end.

Referenced by taskpool_taskprocessor_dtor().

◆ taskpool_taskprocessors_cleanup()

static void taskpool_taskprocessors_cleanup ( struct taskpool_taskprocessors taskprocessors)
static

Definition at line 220 of file taskpool.c.

221{
222 /* Access/manipulation of taskprocessors is done with the lock held, and
223 * with a check of the shutdown flag done. This means that outside of holding
224 * the lock we can safely muck with it. Pushing to the taskprocessor is done
225 * outside of the lock, but with a reference to the taskprocessor held.
226 */
228 AST_VECTOR_FREE(&taskprocessors->taskprocessors);
229}
#define AST_VECTOR_FREE(vec)
Deallocates this vector.
Definition vector.h:185

References ao2_cleanup, AST_VECTOR_CALLBACK_VOID, AST_VECTOR_FREE, and taskpool_taskprocessors::taskprocessors.

Referenced by ast_taskpool_shutdown().

◆ taskpool_taskprocessors_init()

static int taskpool_taskprocessors_init ( struct taskpool_taskprocessors taskprocessors,
unsigned int  size 
)
static

Definition at line 207 of file taskpool.c.

208{
209 if (AST_VECTOR_INIT(&taskprocessors->taskprocessors, size)) {
210 return -1;
211 }
212
213 return 0;
214}
#define AST_VECTOR_INIT(vec, size)
Initialize a vector.
Definition vector.h:124

References AST_VECTOR_INIT, and taskpool_taskprocessors::taskprocessors.

Referenced by ast_taskpool_create().

Variable Documentation

◆ sched

struct ast_sched_context* sched
static

Scheduler used for dynamic pool shrinking.

Definition at line 97 of file taskpool.c.

◆ serializer_tps_listener_callbacks

struct ast_taskprocessor_listener_callbacks serializer_tps_listener_callbacks
static
Initial value:
= {
.task_pushed = serializer_task_pushed,
.start = serializer_start,
.shutdown = serializer_shutdown,
}
static void serializer_task_pushed(struct ast_taskprocessor_listener *listener, int was_empty)
Definition taskpool.c:808
static void serializer_shutdown(struct ast_taskprocessor_listener *listener)
Definition taskpool.c:826
static int serializer_start(struct ast_taskprocessor_listener *listener)
Definition taskpool.c:820

Definition at line 836 of file taskpool.c.

836 {
837 .task_pushed = serializer_task_pushed,
838 .start = serializer_start,
839 .shutdown = serializer_shutdown,
840};

Referenced by ast_taskpool_serializer_group().