-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy paththread_pool.c
More file actions
executable file
·122 lines (114 loc) · 2.92 KB
/
Copy paththread_pool.c
File metadata and controls
executable file
·122 lines (114 loc) · 2.92 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
//
// thread_pool.cpp
// kdlm_v_bufs
//
// Created by yx on 2022/12/20.
//
#include "thread_pool.h"
#include <unistd.h>
#include <stdlib.h>
#include <time.h>
#include <sys/time.h>
#include "time_util.h"
#include "mem.h"
static void* thread_routine(void* arg)
{
kthread_pool_t* pool = (kthread_pool_t*)arg;
struct work_struct* wk;
while (1) {
pthread_mutex_lock(&pool->q_lock);
while (list_empty(&pool->tasks) && !pool->shutdown) {
pthread_cond_wait(&pool->q_cond, &pool->q_lock);
}
if (pool->shutdown) {
pthread_mutex_unlock(&pool->q_lock);
pthread_exit(NULL);
}
wk = list_first_entry(&pool->tasks, struct work_struct, entry);
list_del(&wk->entry);
pthread_mutex_unlock(&pool->q_lock);
wk->func(wk);
}
return NULL;
}
kthread_pool_t* create_pool(int32_t max_thr)
{
int i;
kthread_pool_t* pool = (kthread_pool_t*)calloc_common(1, sizeof(kthread_pool_t));
if (!pool) {
return NULL;
}
pool->thread_count_max = max_thr;
INIT_LIST_HEAD(&pool->tasks);
if (pthread_mutex_init(&pool->q_lock, NULL)) {
free_common(pool);
return NULL;
}
if (pthread_cond_init(&pool->q_cond, NULL)) {
free_common(pool);
return NULL;
}
pool->pt = (pthread_t*)calloc_common(max_thr, sizeof(pthread_t));
if (!pool->pt) {
free_common(pool);
return NULL;
}
for (i = 0; i<max_thr; i++) {
if (pthread_create(&pool->pt[i], NULL, thread_routine, pool)) {
free_common(pool);
free_common(pool->pt);
return NULL;
}
}
return pool;
}
int kthread_pool_add(kthread_pool_t* pool, struct work_struct* wk)
{
if (!pool || !wk) {
return 0;
}
pthread_mutex_lock(&pool->q_lock);
list_add_tail(&wk->entry, &pool->tasks);
pthread_cond_signal(&pool->q_cond);
pthread_mutex_unlock(&pool->q_lock);
return 0;
}
void kthread_pool_destroy(kthread_pool_t* pool)
{
int i;
if (!pool)
return;
if (pool->shutdown)
return;
pool->shutdown = 1;
pthread_mutex_lock(&pool->q_lock);
pthread_cond_broadcast(&pool->q_cond);
pthread_mutex_unlock(&pool->q_lock);
for (i = 0; i<pool->thread_count_max; i++) {
pthread_join(pool->pt[i], NULL);
}
free_common(pool->pt);
INIT_LIST_HEAD(&pool->tasks);
pthread_mutex_destroy(&pool->q_lock);
pthread_cond_destroy(&pool->q_cond);
}
void kthread_pool_flush(kthread_pool_t* pool)
{
if (!pool) {
return;
}
if (pool->shutdown) {
return;
}
while (1) {
pthread_mutex_lock(&pool->q_lock);
if (list_empty(&pool->tasks)) {
pthread_mutex_unlock(&pool->q_lock);
break;
} else {
pthread_cond_broadcast(&pool->q_cond);
}
pthread_mutex_unlock(&pool->q_lock);
}
exactly_usleep(10);
}